diff --git a/fdbclient/BlobGranuleCommon.h b/fdbclient/BlobGranuleCommon.h index 1c24dcb378..c73e76be07 100644 --- a/fdbclient/BlobGranuleCommon.h +++ b/fdbclient/BlobGranuleCommon.h @@ -108,4 +108,5 @@ struct BlobGranuleChunkRef { } }; +enum BlobGranuleSplitState { Unknown = 0, Started = 1, Assigned = 2, Done = 3 }; #endif \ No newline at end of file diff --git a/fdbclient/BlobGranuleReader.actor.cpp b/fdbclient/BlobGranuleReader.actor.cpp index 50c50990d8..626c3d1558 100644 --- a/fdbclient/BlobGranuleReader.actor.cpp +++ b/fdbclient/BlobGranuleReader.actor.cpp @@ -179,6 +179,7 @@ static void applyDelta(std::map* dataMap, Arena& ar, KeyRangeR if (m.param1 < keyRange.begin || m.param1 >= keyRange.end) { return; } + // TODO: we don't need atomics here since eager reads handles it std::map::iterator it = dataMap->find(m.param1); if (m.type != MutationRef::SetValue) { Optional oldVal; @@ -242,15 +243,24 @@ static void applyDeltas(std::map* dataMap, Arena& arena, GranuleDeltas deltas, KeyRangeRef keyRange, - Version readVersion) { + Version readVersion, + Version* lastFileEndVersion) { + if (!deltas.empty()) { + // check that consecutive delta file versions are disjoint + ASSERT(*lastFileEndVersion < deltas.front().version); + } for (MutationsAndVersionRef& delta : deltas) { if (delta.version > readVersion) { - break; + *lastFileEndVersion = readVersion; + return; } for (auto& m : delta.mutations) { applyDelta(dataMap, arena, keyRange, m); } } + if (!deltas.empty()) { + *lastFileEndVersion = deltas.back().version; + } } // TODO: improve the interface of this function so that it doesn't need @@ -272,6 +282,7 @@ ACTOR Future readBlobGranule(BlobGranuleChunkRef chunk, try { state std::map dataMap; + state Version lastFileEndVersion = invalidVersion; Future readSnapshotFuture; if (chunk.snapshotFile.present()) { @@ -301,13 +312,13 @@ ACTOR Future readBlobGranule(BlobGranuleChunkRef chunk, for (Future> deltaFuture : readDeltaFutures) { Standalone result = wait(deltaFuture); arena.dependsOn(result.arena()); - applyDeltas(&dataMap, arena, result, keyRange, readVersion); + applyDeltas(&dataMap, arena, result, keyRange, readVersion, &lastFileEndVersion); wait(yield()); } if (BG_READ_DEBUG) { printf("Applying %d memory deltas\n", chunk.newDeltas.size()); } - applyDeltas(&dataMap, arena, chunk.newDeltas, keyRange, readVersion); + applyDeltas(&dataMap, arena, chunk.newDeltas, keyRange, readVersion, &lastFileEndVersion); wait(yield()); RangeResult ret; diff --git a/fdbclient/BlobWorkerInterface.h b/fdbclient/BlobWorkerInterface.h index f170893992..11d75a48fa 100644 --- a/fdbclient/BlobWorkerInterface.h +++ b/fdbclient/BlobWorkerInterface.h @@ -33,6 +33,8 @@ struct BlobWorkerInterface { RequestStream> waitFailure; RequestStream blobGranuleFileRequest; RequestStream assignBlobRangeRequest; + RequestStream revokeBlobRangeRequest; + RequestStream granuleStatusStreamRequest; struct LocalityData locality; UID myId; @@ -95,23 +97,88 @@ struct AssignBlobRangeReply { } }; -struct AssignBlobRangeRequest { +struct RevokeBlobRangeRequest { constexpr static FileIdentifier file_identifier = 4844288; Arena arena; KeyRangeRef keyRange; int64_t managerEpoch; int64_t managerSeqno; - bool isAssign; // true if assignment, false if revoke + bool dispose; + ReplyPromise reply; + + RevokeBlobRangeRequest() {} + + template + void serialize(Ar& ar) { + serializer(ar, keyRange, managerEpoch, managerSeqno, dispose, reply, arena); + } +}; + +struct AssignBlobRangeRequest { + constexpr static FileIdentifier file_identifier = 905381; + Arena arena; + KeyRangeRef keyRange; + int64_t managerEpoch; + int64_t managerSeqno; + // If continueAssignment is true, this is just to instruct the worker that it still owns the range, so it should + // re-snapshot it and continue. If continueAssignment is false and previousGranules is empty, this is either the + // initial assignment to construct a previously non-existent granule, or a reassignment. Depending on what state + // exists for the granule currently, the worker will either start a new granule, or just pick up from where the + // previous worker left off. + + // For a split or merge, continueAssignment==false. + // For a split, previousGranules will contain one granule that contains keyRange. For a merge, previousGranules will + // contain two or more granules, the union of which will be keyRange. + bool continueAssignment; + VectorRef previousGranules; // only set if there is a granule boundary change + ReplyPromise reply; AssignBlobRangeRequest() {} template void serialize(Ar& ar) { - serializer(ar, keyRange, managerEpoch, managerSeqno, isAssign, reply, arena); + serializer(ar, keyRange, managerEpoch, managerSeqno, continueAssignment, previousGranules, reply, arena); } }; -// TODO once this +// reply per granule +// TODO: could eventually add other types of metrics to report back to the manager here +struct GranuleStatusReply : public ReplyPromiseStreamReply { + constexpr static FileIdentifier file_identifier = 7563104; + + KeyRange granuleRange; + bool doSplit; + int64_t epoch; + int64_t seqno; + + GranuleStatusReply() {} + explicit GranuleStatusReply(KeyRange range, bool doSplit, int64_t epoch, int64_t seqno) + : granuleRange(range), doSplit(doSplit), epoch(epoch), seqno(seqno) {} + + int expectedSize() const { return sizeof(GranuleStatusReply) + granuleRange.expectedSize(); } + + template + void serialize(Ar& ar) { + serializer(ar, ReplyPromiseStreamReply::acknowledgeToken, granuleRange, doSplit, epoch, seqno); + } +}; + +// manager makes one request per worker, it sends all range updates through this stream +struct GranuleStatusStreamRequest { + constexpr static FileIdentifier file_identifier = 2289677; + + int64_t managerEpoch; + + ReplyPromiseStream reply; + + GranuleStatusStreamRequest() {} + explicit GranuleStatusStreamRequest(int64_t managerEpoch) : managerEpoch(managerEpoch) {} + + template + void serialize(Ar& ar) { + serializer(ar, managerEpoch, reply); + } +}; #endif \ No newline at end of file diff --git a/fdbclient/NativeAPI.actor.cpp b/fdbclient/NativeAPI.actor.cpp index d79177d9e6..c8eb0da8bb 100644 --- a/fdbclient/NativeAPI.actor.cpp +++ b/fdbclient/NativeAPI.actor.cpp @@ -6581,7 +6581,7 @@ Future>> DatabaseContext::getRangeF ACTOR Future getRangeFeedStreamActor(Reference db, PromiseStream>> results, - StringRef rangeID, + Key rangeID, Version begin, Version end, KeyRange range) { diff --git a/fdbclient/ServerKnobs.cpp b/fdbclient/ServerKnobs.cpp index 234a2a1a85..20908426c8 100644 --- a/fdbclient/ServerKnobs.cpp +++ b/fdbclient/ServerKnobs.cpp @@ -752,11 +752,14 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi // Blob granlues init( BG_URL, "" ); // TODO CHANGE BACK - init( BG_SNAPSHOT_FILE_TARGET_BYTES, 10000000 ); - // init( BG_SNAPSHOT_FILE_TARGET_BYTES, 1000000 ); + // init( BG_SNAPSHOT_FILE_TARGET_BYTES, 10000000 ); + init( BG_SNAPSHOT_FILE_TARGET_BYTES, 1000000 ); init( BG_DELTA_BYTES_BEFORE_COMPACT, BG_SNAPSHOT_FILE_TARGET_BYTES/2 ); init( BG_DELTA_FILE_TARGET_BYTES, BG_DELTA_BYTES_BEFORE_COMPACT/10 ); + // TODO should discuss proper value for this + init( BLOB_WORKER_TIMEOUT, 10.0 ); if( randomize && BUGGIFY ) BLOB_WORKER_TIMEOUT = 1.0; + // clang-format on if (clientKnobs) { diff --git a/fdbclient/ServerKnobs.h b/fdbclient/ServerKnobs.h index 529e1cdd86..83d09dd685 100644 --- a/fdbclient/ServerKnobs.h +++ b/fdbclient/ServerKnobs.h @@ -704,6 +704,8 @@ public: int BG_DELTA_FILE_TARGET_BYTES; int BG_DELTA_BYTES_BEFORE_COMPACT; + double BLOB_WORKER_TIMEOUT; // Blob Manager's reaction time to a blob worker failure + ServerKnobs(Randomize, ClientKnobs*, IsSimulated); void initialize(Randomize, ClientKnobs*, IsSimulated); }; diff --git a/fdbclient/SystemData.cpp b/fdbclient/SystemData.cpp index bdfb1d1380..07e4be6d4a 100644 --- a/fdbclient/SystemData.cpp +++ b/fdbclient/SystemData.cpp @@ -1104,6 +1104,7 @@ int64_t decodeBlobManagerEpochValue(ValueRef const& value) { const KeyRangeRef blobGranuleFileKeys(LiteralStringRef("\xff\x02/bgf/"), LiteralStringRef("\xff\x02/bgf0")); const KeyRangeRef blobGranuleMappingKeys(LiteralStringRef("\xff\x02/bgm/"), LiteralStringRef("\xff\x02/bgm0")); const KeyRangeRef blobGranuleLockKeys(LiteralStringRef("\xff\x02/bgl/"), LiteralStringRef("\xff\x02/bgl0")); +const KeyRangeRef blobGranuleSplitKeys(LiteralStringRef("\xff\x02/bgs/"), LiteralStringRef("\xff\x02/bgs0")); const Value blobGranuleMappingValueFor(UID const& workerID) { BinaryWriter wr(Unversioned()); @@ -1118,19 +1119,35 @@ UID decodeBlobGranuleMappingValue(ValueRef const& value) { return workerID; } -const Value blobGranuleLockValueFor(int64_t epoch, int64_t seqno) { +const Value blobGranuleLockValueFor(int64_t epoch, int64_t seqno, UID changeFeedId) { BinaryWriter wr(Unversioned()); wr << epoch; wr << seqno; + wr << changeFeedId; return wr.toValue(); } -std::pair decodeBlobGranuleLockValue(const ValueRef& value) { +std::tuple decodeBlobGranuleLockValue(const ValueRef& value) { int64_t epoch, seqno; + UID changeFeedId; BinaryReader reader(value, Unversioned()); reader >> epoch; reader >> seqno; - return std::pair(epoch, seqno); + reader >> changeFeedId; + return std::make_tuple(epoch, seqno, changeFeedId); +} + +const Value blobGranuleSplitValueFor(BlobGranuleSplitState st) { + BinaryWriter wr(Unversioned()); + wr << st; + return wr.toValue(); +} + +BlobGranuleSplitState decodeBlobGranuleSplitValue(const ValueRef& value) { + BlobGranuleSplitState st; + BinaryReader reader(value, Unversioned()); + reader >> st; + return st; } const KeyRangeRef blobWorkerListKeys(LiteralStringRef("\xff\x02/bwList/"), LiteralStringRef("\xff\x02/bwList0")); diff --git a/fdbclient/SystemData.h b/fdbclient/SystemData.h index df405e9650..8e358c4869 100644 --- a/fdbclient/SystemData.h +++ b/fdbclient/SystemData.h @@ -25,7 +25,7 @@ // Functions and constants documenting the organization of the reserved keyspace in the database beginning with "\xFF" #include "fdbclient/FDBTypes.h" -#include "fdbclient/BlobWorkerInterface.h" // TODO move the functions that depend on this out of here and into BlobWorkerInterface.h +#include "fdbclient/BlobWorkerInterface.h" // TODO move the functions that depend on this out of here and into BlobWorkerInterface.h to remove this depdendency #include "fdbclient/StorageServerInterface.h" // Don't warn on constants being defined in this file. @@ -535,16 +535,23 @@ extern const KeyRangeRef blobGranuleFileKeys; // \xff/bgm/[[begin]] = [[BlobWorkerUID]] extern const KeyRangeRef blobGranuleMappingKeys; -// \xff/bgl/(begin,end) = (epoch, seqno) +// \xff/bgl/(begin,end) = (epoch, seqno, changefeed id) extern const KeyRangeRef blobGranuleLockKeys; -const Value blobGranuleLockValueFor(int64_t epochNum, int64_t sequenceNum); -std::pair decodeBlobGranuleLockValue(ValueRef const& value); +// \xff/bgs/(oldbegin,oldend,newbegin) = state +extern const KeyRangeRef blobGranuleSplitKeys; const Value blobGranuleMappingValueFor(UID const& workerID); UID decodeBlobGranuleMappingValue(ValueRef const& value); -// \xff/blobWorkerList/[[BlobWorkerID]] = [[BlobWorkerInterface]] +const Value blobGranuleLockValueFor(int64_t epochNum, int64_t sequenceNum, UID changeFeedId); +// FIXME: maybe just define a struct? +std::tuple decodeBlobGranuleLockValue(ValueRef const& value); + +const Value blobGranuleSplitValueFor(BlobGranuleSplitState st); +BlobGranuleSplitState decodeBlobGranuleSplitValue(ValueRef const& value); + +// \xff/bwl/[[BlobWorkerID]] = [[BlobWorkerInterface]] extern const KeyRangeRef blobWorkerListKeys; const Key blobWorkerListKeyFor(UID workerID); diff --git a/fdbserver/BlobManager.actor.cpp b/fdbserver/BlobManager.actor.cpp index a807097d56..f002cd5acb 100644 --- a/fdbserver/BlobManager.actor.cpp +++ b/fdbserver/BlobManager.actor.cpp @@ -25,9 +25,11 @@ #include "fdbclient/KeyRangeMap.h" #include "fdbclient/ReadYourWrites.h" #include "fdbclient/SystemData.h" +#include "fdbclient/Tuple.h" #include "fdbserver/BlobManagerInterface.h" #include "fdbserver/BlobWorker.actor.h" #include "fdbserver/Knobs.h" +#include "fdbserver/WaitFailure.h" #include "flow/IRandom.h" #include "flow/UnitTest.h" #include "flow/actorcompiler.h" // has to be last include @@ -160,23 +162,37 @@ void getRanges(std::vector>& results, KeyRangeMap previousRanges; - RangeAssignment() {} - explicit RangeAssignment(KeyRange keyRange, bool isAssign) : keyRange(keyRange), isAssign(isAssign) {} + RangeAssignmentData() : continueAssignment(false) {} + RangeAssignmentData(bool continueAssignment, std::vector previousRanges) + : continueAssignment(continueAssignment), previousRanges(previousRanges) {} +}; + +struct RangeRevokeData { + bool dispose; + + RangeRevokeData() {} + RangeRevokeData(bool dispose) : dispose(dispose) {} +}; + +struct RangeAssignment { + bool isAssign; + KeyRange keyRange; + Optional worker; + + // I tried doing this with a union and it was just kind of messy + Optional assign; + Optional revoke; }; // TODO: track worker's reads/writes eventually struct BlobWorkerStats { int numGranulesAssigned; - BlobWorkerStats(int numGranulesAssigned=0): numGranulesAssigned(numGranulesAssigned) {} + BlobWorkerStats(int numGranulesAssigned = 0) : numGranulesAssigned(numGranulesAssigned) {} }; struct BlobManagerData { @@ -185,6 +201,7 @@ struct BlobManagerData { std::unordered_map workersById; std::unordered_map workerStats; // mapping between workerID -> workerStats + std::unordered_map> workerMonitors; KeyRangeMap workerAssignments; KeyRangeMap knownBlobRanges; @@ -222,7 +239,7 @@ ACTOR Future nukeBlobWorkerData(BlobManagerData* bmData) { } } -ACTOR Future>> splitNewRange(Reference tr, KeyRange range) { +ACTOR Future>> splitRange(Reference tr, KeyRange range) { // TODO is it better to just pass empty metrics to estimated? // TODO handle errors here by pulling out into its own transaction instead of the main loop's transaction, and // retrying @@ -264,8 +281,8 @@ ACTOR Future>> splitNewRange(Reference eligibleWorkers; - - for (auto const &worker : bmData->workerStats) { + + for (auto const& worker : bmData->workerStats) { UID currId = worker.first; int granulesAssigned = worker.second.numGranulesAssigned; @@ -282,33 +299,57 @@ static UID pickWorkerForAssign(BlobManagerData* bmData) { ASSERT(eligibleWorkers.size() > 0); int idx = deterministicRandom()->randomInt(0, eligibleWorkers.size()); if (BM_DEBUG) { - printf("picked worker %s, which has a minimal number (%d) of granules assigned\n", - eligibleWorkers[idx].toString().c_str(), minGranulesAssigned); + printf("picked worker %s, which has a minimal number (%d) of granules assigned\n", + eligibleWorkers[idx].toString().c_str(), + minGranulesAssigned); } return eligibleWorkers[idx]; } ACTOR Future doRangeAssignment(BlobManagerData* bmData, RangeAssignment assignment, UID workerID, int64_t seqNo) { - AssignBlobRangeRequest req; - req.keyRange = - KeyRangeRef(StringRef(req.arena, assignment.keyRange.begin), StringRef(req.arena, assignment.keyRange.end)); - req.managerEpoch = bmData->epoch; - req.managerSeqno = seqNo; - req.isAssign = assignment.isAssign; if (BM_DEBUG) { printf("BM %s %s range [%s - %s) @ (%lld, %lld)\n", workerID.toString().c_str(), - req.isAssign ? "assigning" : "revoking", - req.keyRange.begin.printable().c_str(), - req.keyRange.end.printable().c_str(), - req.managerEpoch, - req.managerSeqno); + assignment.isAssign ? "assigning" : "revoking", + assignment.keyRange.begin.printable().c_str(), + assignment.keyRange.end.printable().c_str(), + bmData->epoch, + seqNo); } try { - AssignBlobRangeReply rep = wait(bmData->workersById[workerID].assignBlobRangeRequest.getReply(req)); + state AssignBlobRangeReply rep; + if (assignment.isAssign) { + ASSERT(assignment.assign.present()); + ASSERT(!assignment.revoke.present()); + + AssignBlobRangeRequest req; + req.keyRange = KeyRangeRef(StringRef(req.arena, assignment.keyRange.begin), + StringRef(req.arena, assignment.keyRange.end)); + req.managerEpoch = bmData->epoch; + req.managerSeqno = seqNo; + req.continueAssignment = assignment.assign.get().continueAssignment; + for (auto& it : assignment.assign.get().previousRanges) { + req.previousGranules.push_back_deep(req.arena, it); + } + AssignBlobRangeReply _rep = wait(bmData->workersById[workerID].assignBlobRangeRequest.getReply(req)); + rep = _rep; + } else { + ASSERT(!assignment.assign.present()); + ASSERT(assignment.revoke.present()); + + RevokeBlobRangeRequest req; + req.keyRange = KeyRangeRef(StringRef(req.arena, assignment.keyRange.begin), + StringRef(req.arena, assignment.keyRange.end)); + req.managerEpoch = bmData->epoch; + req.managerSeqno = seqNo; + req.dispose = assignment.revoke.get().dispose; + + AssignBlobRangeReply _rep = wait(bmData->workersById[workerID].revokeBlobRangeRequest.getReply(req)); + rep = _rep; + } if (!rep.epochOk) { if (BM_DEBUG) { printf("BM heard from BW that there is a new manager with higher epoch\n"); @@ -327,13 +368,37 @@ ACTOR Future doRangeAssignment(BlobManagerData* bmData, RangeAssignment as assignment.keyRange.end.printable().c_str()); } // re-send revoke to queue to handle range being un-assigned from that worker before the new one - bmData->rangesToAssign.send(RangeAssignment(assignment.keyRange, false)); + RangeAssignment revokeOld; + revokeOld.isAssign = false; + revokeOld.worker = workerID; + revokeOld.keyRange = assignment.keyRange; + revokeOld.revoke = RangeRevokeData(false); + bmData->rangesToAssign.send(revokeOld); + + // send assignment back to queue as is, clearing designated worker if present + assignment.worker.reset(); bmData->rangesToAssign.send(assignment); // FIXME: improvement would be to add history of failed workers to assignment so it can try other ones first - } else if (BM_DEBUG) { - printf("BM got error revoking range [%s - %s) from worker %s, ignoring\n", - assignment.keyRange.begin.printable().c_str(), - assignment.keyRange.end.printable().c_str()); + } else { + if (BM_DEBUG) { + printf("BM got error revoking range [%s - %s) from worker %s", + assignment.keyRange.begin.printable().c_str(), + assignment.keyRange.end.printable().c_str()); + } + + if (assignment.revoke.get().dispose) { + if (BM_DEBUG) { + printf(", retrying for dispose\n"); + } + // send assignment back to queue as is, clearing designated worker if present + assignment.worker.reset(); + bmData->rangesToAssign.send(assignment); + // + } else { + if (BM_DEBUG) { + printf(", ignoring\n"); + } + } } } return Void(); @@ -354,12 +419,17 @@ ACTOR Future rangeAssigner(BlobManagerData* bmData) { auto currentAssignments = bmData->workerAssignments.intersectingRanges(assignment.keyRange); int count = 0; for (auto& it : currentAssignments) { - ASSERT(it.value() == UID()); + if (assignment.assign.get().continueAssignment) { + ASSERT(assignment.worker.present()); + ASSERT(it.value() == assignment.worker.get()); + } else { + ASSERT(it.value() == UID()); + } count++; } ASSERT(count == 1); - workerId = pickWorkerForAssign(bmData); + workerId = assignment.worker.present() ? assignment.worker.get() : pickWorkerForAssign(bmData); bmData->workerAssignments.insert(assignment.keyRange, workerId); bmData->workerStats[workerId].numGranulesAssigned += 1; @@ -377,7 +447,8 @@ ACTOR Future rangeAssigner(BlobManagerData* bmData) { // It is fine for multiple disjoint sub-ranges to have the same sequence number since they were part of // the same logical change bmData->workerStats[it.value()].numGranulesAssigned -= 1; - addActor.send(doRangeAssignment(bmData, assignment, it.value(), seqNo)); + if (!assignment.worker.present() || assignment.worker.get() == it.value()) + addActor.send(doRangeAssignment(bmData, assignment, it.value(), seqNo)); } bmData->workerAssignments.insert(assignment.keyRange, UID()); @@ -385,14 +456,37 @@ ACTOR Future rangeAssigner(BlobManagerData* bmData) { } } +ACTOR Future checkManagerLock(Reference tr, BlobManagerData* bmData) { + Optional currentLockValue = wait(tr->get(blobManagerEpochKey)); + ASSERT(currentLockValue.present()); + int64_t currentEpoch = decodeBlobManagerEpochValue(currentLockValue.get()); + if (currentEpoch != bmData->epoch) { + ASSERT(currentEpoch > bmData->epoch); + + printf("BM %s found new epoch %d > %d in lock check\n", + bmData->id.toString().c_str(), + currentEpoch, + bmData->epoch); + if (bmData->iAmReplaced.canBeSet()) { + bmData->iAmReplaced.send(Void()); + } + + // TODO different error? + throw granule_assignment_conflict(); + } + tr->addReadConflictRange(singleKeyRange(blobManagerEpochKey)); + + return Void(); +} + // TODO eventually CC should probably do this and pass it as part of recruitment? ACTOR Future acquireManagerLock(BlobManagerData* bmData) { state Reference tr = makeReference(bmData->db); - tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS); - tr->setOption(FDBTransactionOptions::PRIORITY_SYSTEM_IMMEDIATE); + loop { + tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS); + tr->setOption(FDBTransactionOptions::PRIORITY_SYSTEM_IMMEDIATE); try { - // TODO verify: this should automatically have a read conflict range for blobManagerEpochKey, right? Optional oldEpoch = wait(tr->get(blobManagerEpochKey)); state int64_t newEpoch; if (oldEpoch.present()) { @@ -414,6 +508,8 @@ ACTOR Future acquireManagerLock(BlobManagerData* bmData) { } } +// FIXME: this does all logic in one transaction. Adding a giant range to an existing database to hybridize would spread +// require doing a ton of storage metrics calls, which we should split up across multiple transactions likely. ACTOR Future monitorClientRanges(BlobManagerData* bmData) { loop { state Reference tr = makeReference(bmData->db); @@ -445,14 +541,18 @@ ACTOR Future monitorClientRanges(BlobManagerData* bmData) { range.end.printable().c_str()); } - bmData->rangesToAssign.send(RangeAssignment(range, false)); + RangeAssignment ra; + ra.isAssign = false; + ra.keyRange = range; + ra.revoke = RangeRevokeData(true); // dispose=true + bmData->rangesToAssign.send(ra); } state std::vector>>> splitFutures; // Divide new ranges up into equal chunks by using SS byte sample for (KeyRangeRef range : rangesToAdd) { // assert that this range contains no currently assigned ranges in this - splitFutures.push_back(splitNewRange(tr, range)); + splitFutures.push_back(splitRange(tr, range)); } for (auto f : splitFutures) { @@ -470,7 +570,11 @@ ACTOR Future monitorClientRanges(BlobManagerData* bmData) { printf(" [%s - %s)\n", range.begin.printable().c_str(), range.end.printable().c_str()); } - bmData->rangesToAssign.send(RangeAssignment(range, true)); + RangeAssignment ra; + ra.isAssign = true; + ra.keyRange = range; + ra.assign = RangeAssignmentData(); // continue=false, no previous granules + bmData->rangesToAssign.send(ra); } } @@ -491,6 +595,213 @@ ACTOR Future monitorClientRanges(BlobManagerData* bmData) { } } +static Key granuleLockKey(KeyRange granuleRange) { + Tuple k; + k.append(granuleRange.begin).append(granuleRange.end); + return k.getDataAsStandalone().withPrefix(blobGranuleLockKeys.begin); +} + +// FIXME: propagate errors here +ACTOR Future maybeSplitRange(BlobManagerData* bmData, UID currentWorkerId, KeyRange range) { + state Reference tr = makeReference(bmData->db); + state Standalone> newRanges; + state int64_t newLockSeqno = -1; + + // first get ranges to split + loop { + try { + // redo split if previous txn try failed to calculate it + if (newRanges.empty()) { + Standalone> _newRanges = wait(splitRange(tr, range)); + newRanges = _newRanges; + } + break; + } catch (Error& e) { + wait(tr->onError(e)); + } + } + + if (newRanges.size() == 2) { + // not large enough to split, just reassign back to worker + if (BM_DEBUG) { + printf("Not splitting existing range [%s - %s). Continuing assignment to %s\n", + range.begin.printable().c_str(), + range.end.printable().c_str(), + currentWorkerId.toString().c_str()); + } + RangeAssignment raContinue; + raContinue.isAssign = true; + raContinue.worker = currentWorkerId; + raContinue.keyRange = range; + raContinue.assign = + RangeAssignmentData(true, std::vector()); // continue, no "previous" range to do handover + bmData->rangesToAssign.send(raContinue); + return Void(); + } + + // Need to split range. Persist intent to split and split metadata to DB BEFORE sending split requests + loop { + try { + tr->reset(); + tr->setOption(FDBTransactionOptions::Option::PRIORITY_SYSTEM_IMMEDIATE); + tr->setOption(FDBTransactionOptions::Option::ACCESS_SYSTEM_KEYS); + ASSERT(newRanges.size() >= 2); + + // make sure we're still manager when this transaction gets committed + wait(checkManagerLock(tr, bmData)); + + // acquire lock for old granule to make sure nobody else modifies it + state Key lockKey = granuleLockKey(range); + Optional lockValue = wait(tr->get(lockKey)); + ASSERT(lockValue.present()); + std::tuple prevGranuleLock = decodeBlobGranuleLockValue(lockValue.get()); + if (std::get<0>(prevGranuleLock) > bmData->epoch) { + printf("BM %s found a higher epoch %d than %d for granule lock of [%s - %s)\n", + bmData->id.toString().c_str(), + std::get<0>(prevGranuleLock), + bmData->epoch, + range.begin.printable().c_str(), + range.end.printable().c_str()); + + if (bmData->iAmReplaced.canBeSet()) { + bmData->iAmReplaced.send(Void()); + } + return Void(); + } + if (newLockSeqno == -1) { + newLockSeqno = bmData->seqNo; + bmData->seqNo++; + ASSERT(newLockSeqno > std::get<1>(prevGranuleLock)); + } else { + // previous transaction could have succeeded but got commit_unknown_result + ASSERT(newLockSeqno >= std::get<1>(prevGranuleLock)); + } + + tr->set(lockKey, blobGranuleLockValueFor(bmData->epoch, newLockSeqno, std::get<2>(prevGranuleLock))); + + // set up split metadata + for (int i = 0; i < newRanges.size() - 1; i++) { + Tuple key; + key.append(range.begin).append(range.end).append(newRanges[i]); + tr->set(key.getDataAsStandalone().withPrefix(blobGranuleSplitKeys.begin), + blobGranuleSplitValueFor(BlobGranuleSplitState::Started)); + + // acquire granule lock so nobody else can make changes to this granule. + } + wait(tr->commit()); + break; + } catch (Error& e) { + wait(tr->onError(e)); + } + } + + if (BM_DEBUG) { + printf("Splitting range [%s - %s) into:\n", range.begin.printable().c_str(), range.end.printable().c_str()); + for (int i = 0; i < newRanges.size() - 1; i++) { + printf(" [%s - %s)\n", newRanges[i].printable().c_str(), newRanges[i + 1].printable().c_str()); + } + } + + // transaction committed, send range assignments + // revoke from current worker + RangeAssignment raRevoke; + raRevoke.isAssign = false; + raRevoke.worker = currentWorkerId; + raRevoke.keyRange = range; + raRevoke.revoke = RangeRevokeData(false); // not a dispose + bmData->rangesToAssign.send(raRevoke); + + std::vector originalRange; + originalRange.push_back(range); + for (int i = 0; i < newRanges.size() - 1; i++) { + // reassign new range and do handover of previous range + RangeAssignment raAssignSplit; + raAssignSplit.isAssign = true; + raAssignSplit.keyRange = KeyRangeRef(newRanges[i], newRanges[i + 1]); + raAssignSplit.assign = RangeAssignmentData(false, originalRange); + // don't care who this range gets assigned to + bmData->rangesToAssign.send(raAssignSplit); + } + + return Void(); +} + +ACTOR Future monitorBlobWorker(BlobManagerData* bmData, BlobWorkerInterface bwInterf) { + try { + state PromiseStream> addActor; + state Future collection = actorCollection(addActor.getFuture()); + state Future waitFailure = waitFailureClient(bwInterf.waitFailure, SERVER_KNOBS->BLOB_WORKER_TIMEOUT); + state ReplyPromiseStream statusStream = + bwInterf.granuleStatusStreamRequest.getReplyStream(GranuleStatusStreamRequest(bmData->epoch)); + state KeyRangeMap> lastSeenSeqno; + + loop choose { + when(wait(waitFailure)) { + // FIXME: actually handle this!! + if (BM_DEBUG) { + printf("BM %lld detected BW %s is dead\n", bmData->epoch, bwInterf.id().toString().c_str()); + } + return Void(); + } + when(GranuleStatusReply _rep = waitNext(statusStream.getFuture())) { + GranuleStatusReply rep = _rep; + if (BM_DEBUG) { + printf("BM %lld got status of [%s - %s) @ (%lld, %lld) from BW %s: %s\n", + bmData->epoch, + rep.granuleRange.begin.printable().c_str(), + rep.granuleRange.end.printable().c_str(), + rep.epoch, + rep.seqno, + bwInterf.id().toString().c_str(), + rep.doSplit ? "split" : ""); + } + if (rep.epoch > bmData->epoch) { + if (BM_DEBUG) { + printf("BM heard from BW that there is a new manager with higher epoch\n"); + } + if (bmData->iAmReplaced.canBeSet()) { + bmData->iAmReplaced.send(Void()); + } + } + + // TODO maybe this won't be true eventually, but right now the only time the blob worker reports back is + // to split the range. + ASSERT(rep.doSplit); + + // FIXME: only evaluate for split if this worker currently owns the granule in this blob manager's + // mapping + + auto lastReqForGranule = lastSeenSeqno.rangeContaining(rep.granuleRange.begin); + if (rep.granuleRange.begin == lastReqForGranule.begin() && + rep.granuleRange.end == lastReqForGranule.end() && rep.epoch == lastReqForGranule.value().first && + rep.seqno == lastReqForGranule.value().second) { + if (BM_DEBUG) { + printf("Manager %lld received repeat status for the same granule [%s - %s) @ %lld, ignoring.", + bmData->epoch, + rep.granuleRange.begin.printable().c_str(), + rep.granuleRange.end.printable().c_str()); + } + } else { + if (BM_DEBUG) { + printf("Manager %lld evaluating [%s - %s) for split\n", + bmData->epoch, + rep.granuleRange.begin.printable().c_str(), + rep.granuleRange.end.printable().c_str()); + } + lastSeenSeqno.insert(rep.granuleRange, std::pair(rep.epoch, rep.seqno)); + addActor.send(maybeSplitRange(bmData, bwInterf.id(), rep.granuleRange)); + } + } + } + } catch (Error& e) { + // FIXME: forward errors somewhere from here + if (BM_DEBUG) { + printf("BM got unexpected error %s monitoring BW %s\n", e.name(), bwInterf.id().toString().c_str()); + } + throw e; + } +} + // TODO this is only for chaos testing right now!! REMOVE LATER ACTOR Future rangeMover(BlobManagerData* bmData) { loop { @@ -509,8 +820,19 @@ ACTOR Future rangeMover(BlobManagerData* bmData) { randomRange.value().toString().c_str()); } - bmData->rangesToAssign.send(RangeAssignment(randomRange.range(), false)); - bmData->rangesToAssign.send(RangeAssignment(randomRange.range(), true)); + RangeAssignment revokeOld; + revokeOld.isAssign = false; + revokeOld.keyRange = randomRange.range(); + revokeOld.worker = randomRange.value(); + revokeOld.revoke = RangeRevokeData(false); + bmData->rangesToAssign.send(revokeOld); + + RangeAssignment assignNew; + assignNew.isAssign = true; + assignNew.keyRange = randomRange.range(); + assignNew.assign = + RangeAssignmentData(false, std::vector()); // not a continue, no boundary change + bmData->rangesToAssign.send(assignNew); break; } } @@ -555,6 +877,7 @@ ACTOR Future blobManager(LocalityData locality, Reference { +struct GranuleFiles { std::deque snapshotFiles; std::deque deltaFiles; +}; + +// TODO needs better name, it's basically just "granule starting state" +struct GranuleChangeFeedInfo { + UID changeFeedId; + Version changeFeedStartVersion; + Version previousDurableVersion; + Optional prevChangeFeedId; + bool doSnapshot; + Optional granuleSplitFrom; + Optional blobFilesToSnapshot; +}; + +// FIXME: the circular dependencies here are getting kind of gross +struct GranuleMetadata; +struct BlobWorkerData; +ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwData, AssignBlobRangeRequest req); +ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, Reference metadata); + +// for a range that is active +struct GranuleMetadata : NonCopyable, ReferenceCounted { + GranuleFiles files; GranuleDeltas currentDeltas; uint64_t bytesInNewDeltaFiles = 0; Version lastWriteVersion = 0; @@ -58,20 +83,42 @@ struct GranuleMetadata : NonCopyable, ReferenceCounted { uint64_t currentDeltaBytes = 0; Arena deltaArena; - int64_t lockEpoch; - int64_t lockSeqno; + int64_t originalEpoch; + int64_t originalSeqno; + int64_t continueEpoch; + int64_t continueSeqno; KeyRange keyRange; - Future assignFuture; + Future assignFuture; Future fileUpdaterFuture; + Promise resumeSnapshot; + + Future start(BlobWorkerData* bwData, AssignBlobRangeRequest req) { + assignFuture = persistAssignWorkerRange(bwData, req); + fileUpdaterFuture = blobGranuleUpdateFiles(bwData, Reference::addRef(this)); + + return success(assignFuture); + } + + void resume() { + ASSERT(resumeSnapshot.canBeSet()); + resumeSnapshot.send(Void()); + } // FIXME: right now there is a dependency because this contains both the actual file/delta data as well as the // metadata (worker futures), so removing this reference from the map doesn't actually cancel the workers. It'd be // better to have this in 2 separate objects, where the granule metadata map has the futures, but the read // queries/file updater/range feed only copy the reference to the file/delta data. - void cancel() { - assignFuture = Never(); - fileUpdaterFuture = Never(); + Future cancel(bool dispose) { + assignFuture.cancel(); + fileUpdaterFuture.cancel(); + + if (dispose) { + // FIXME: implement dispose! + return delay(0.1); + } else { + return Future(Void()); + } } }; @@ -95,51 +142,75 @@ struct BlobWorkerData { LocalityData locality; int64_t currentManagerEpoch = -1; + ReplyPromiseStream currentManagerStatusStream; + // FIXME: refactor out the parts of this that are just for interacting with blob stores from the backup business // logic Reference bstore; - // Reference bstore; - // std::string bucket; - KeyRangeMap granuleMetadata; BlobWorkerData(UID id, Database db) : id(id), db(db), stats(id, SERVER_KNOBS->WORKER_LOGGING_INTERVAL) {} ~BlobWorkerData() { printf("Destroying blob worker data for %s\n", id.toString().c_str()); } + + bool managerEpochOk(int64_t epoch) { + if (epoch < currentManagerEpoch) { + if (BW_DEBUG) { + printf("BW %s got request from old epoch %lld, notifying manager it is out of date\n", + id.toString().c_str(), + epoch); + } + return false; + } else { + if (epoch > currentManagerEpoch) { + currentManagerEpoch = epoch; + if (BW_DEBUG) { + printf("BW %s found new manager epoch %lld\n", id.toString().c_str(), currentManagerEpoch); + } + } + + return true; + } + } }; // returns true if we can acquire it -static void acquireGranuleLock(int64_t epoch, int64_t seqno, std::pair prevOwner) { +static void acquireGranuleLock(int64_t epoch, int64_t seqno, int64_t prevOwnerEpoch, int64_t prevOwnerSeqno) { // returns true if our lock (E, S) >= (Eprev, Sprev) - if (epoch < prevOwner.first || (epoch == prevOwner.first && seqno < prevOwner.second)) { + if (epoch < prevOwnerEpoch || (epoch == prevOwnerEpoch && seqno < prevOwnerSeqno)) { if (BW_DEBUG) { printf("Lock acquire check failed. Proposed (%lld, %lld) < previous (%lld, %lld)\n", epoch, seqno, - prevOwner.first, - prevOwner.second); + prevOwnerEpoch, + prevOwnerSeqno); } throw granule_assignment_conflict(); } } -static void checkGranuleLock(int64_t epoch, int64_t seqno, std::pair currentOwner) { +static void checkGranuleLock(int64_t epoch, int64_t seqno, int64_t ownerEpoch, int64_t ownerSeqno) { // sanity check - lock value should never go backwards because of acquireGranuleLock - ASSERT(epoch <= currentOwner.first); - ASSERT(epoch < currentOwner.first || (epoch == currentOwner.first && seqno <= currentOwner.second)); + /* + printf( + "Checking granule lock: \n mine: (%lld, %lld)\n owner: (%lld, %lld)\n", epoch, seqno, ownerEpoch, ownerSeqno); + */ + ASSERT(epoch <= ownerEpoch); + ASSERT(epoch < ownerEpoch || (epoch == ownerEpoch && seqno <= ownerSeqno)); // returns true if we still own the lock, false if someone else does - if (epoch != currentOwner.first || seqno != currentOwner.second) { + if (epoch != ownerEpoch || seqno != ownerSeqno) { if (BW_DEBUG) { printf("Lock assignment check failed. Expected (%lld, %lld), got (%lld, %lld)\n", epoch, seqno, - currentOwner.first, - currentOwner.second); + ownerEpoch, + ownerSeqno); } throw granule_assignment_conflict(); } } +// TODO this is duplicated with blob manager: fix? static Key granuleLockKey(KeyRange granuleRange) { Tuple k; k.append(granuleRange.begin).append(granuleRange.end); @@ -154,8 +225,8 @@ ACTOR Future readAndCheckGranuleLock(Reference Optional lockValue = wait(tr->get(lockKey)); ASSERT(lockValue.present()); - std::pair currentOwner = decodeBlobGranuleLockValue(lockValue.get()); - checkGranuleLock(epoch, seqno, currentOwner); + std::tuple currentOwner = decodeBlobGranuleLockValue(lockValue.get()); + checkGranuleLock(epoch, seqno, std::get<0>(currentOwner), std::get<1>(currentOwner)); // if we still own the lock, add a conflict range in case anybody else takes it over while we add this file tr->addReadConflictRange(singleKeyRange(lockKey)); @@ -163,6 +234,186 @@ ACTOR Future readAndCheckGranuleLock(Reference return Void(); } +ACTOR Future loadPreviousFiles(Transaction* tr, KeyRange keyRange) { + // read everything from previous granule of snapshot and delta files + Tuple prevFilesStartTuple; + prevFilesStartTuple.append(keyRange.begin).append(keyRange.end); + Tuple prevFilesEndTuple; + prevFilesEndTuple.append(keyRange.begin).append(keyAfter(keyRange.end)); + Key prevFilesStartKey = prevFilesStartTuple.getDataAsStandalone().withPrefix(blobGranuleFileKeys.begin); + Key prevFilesEndKey = prevFilesEndTuple.getDataAsStandalone().withPrefix(blobGranuleFileKeys.begin); + + state KeyRange currentRange = KeyRangeRef(prevFilesStartKey, prevFilesEndKey); + + state GranuleFiles files; + + loop { + RangeResult res = wait(tr->getRange(currentRange, 1000)); + for (auto& it : res) { + Tuple fileKey = Tuple::unpack(it.key.removePrefix(blobGranuleFileKeys.begin)); + Tuple fileValue = Tuple::unpack(it.value); + + ASSERT(fileKey.size() == 4); + ASSERT(fileValue.size() == 3); + + ASSERT(fileKey.getString(0) == keyRange.begin); + ASSERT(fileKey.getString(1) == keyRange.end); + + std::string fileType = fileKey.getString(2).toString(); + ASSERT(fileType == LiteralStringRef("S") || fileType == LiteralStringRef("D")); + + BlobFileIndex idx( + fileKey.getInt(3), fileValue.getString(0).toString(), fileValue.getInt(1), fileValue.getInt(2)); + if (fileType == LiteralStringRef("S")) { + ASSERT(files.snapshotFiles.empty() || files.snapshotFiles.back().version < idx.version); + files.snapshotFiles.push_back(idx); + } else { + ASSERT(files.deltaFiles.empty() || files.deltaFiles.back().version < idx.version); + files.deltaFiles.push_back(idx); + } + } + if (res.more) { + currentRange = KeyRangeRef(keyAfter(res.back().key), currentRange.end); + } else { + break; + } + } + printf("Loaded %d snapshot and %d delta previous files for [%s - %s)\n", + files.snapshotFiles.size(), + files.deltaFiles.size(), + keyRange.begin.printable().c_str(), + keyRange.end.printable().c_str()); + return files; +} + +// To cleanup of the old change feed for the old granule range, all new sub-granules split from the old range must +// update shared state to coordinate when it is safe to clean up the old change feed. +// his goes through 3 phases for each new sub-granule: +// 1. Starting - the blob manager writes all sub-granules with this state as a durable intent to split the range +// 2. Assigned - a worker that is assigned a sub-granule updates that granule's state here. This means that the worker +// has started a new change feed for the new sub-granule, but still needs to consume from the old change feed. +// 3. Done - the worker that is assigned this sub-granule has persisted all of the data from its part of the old change +// feed in delta files. From this granule's perspective, it is safe to clean up the old change feed. + +// Once all sub-granules have reached step 2 (Assigned), the change feed can be safely "stopped" - it needs to continue +// to serve the mutations it has seen so far, but will not need any new mutations after this version. +// The last sub-granule to reach this step is responsible for commiting the change feed stop as part of its +// transaction. Because this change feed stops commits in the same transaction as the worker's new change feed start, +// it is guaranteed that no versions are missed between the old and new change feed. +// +// Once all sub-granules have reached step 3 (Done), the change feed can be safely destroyed, as all of the mutations in +// the old change feed are guaranteed to be persisted in delta files. The last sub-granule to reach this step is +// responsible for committing the change feed destroy, and for cleaning up the split state for all sub-granules as part +// of its transaction. + +ACTOR Future updateGranuleSplitState(Transaction* tr, + KeyRange previousGranule, + KeyRange currentGranule, + UID prevChangeFeedId, + BlobGranuleSplitState newState) { + // read all splitting state for previousGranule. If it doesn't exist, newState must == DONE + Tuple splitStateStartTuple; + splitStateStartTuple.append(previousGranule.begin).append(previousGranule.end); + Tuple splitStateEndTuple; + splitStateEndTuple.append(previousGranule.begin).append(keyAfter(previousGranule.end)); + + Key splitStateStartKey = splitStateStartTuple.getDataAsStandalone().withPrefix(blobGranuleSplitKeys.begin); + Key splitStateEndKey = splitStateEndTuple.getDataAsStandalone().withPrefix(blobGranuleSplitKeys.begin); + + state KeyRange currentRange = KeyRangeRef(splitStateStartKey, splitStateEndKey); + + RangeResult totalState = wait(tr->getRange(currentRange, 10)); + // TODO is this explicit conflit range necessary with the above read? + tr->addWriteConflictRange(currentRange); + ASSERT(!totalState.more); + + if (totalState.empty()) { + ASSERT(newState == BlobGranuleSplitState::Done); + printf("Found empty split state for previous granule [%s - %s)\n", + previousGranule.begin.printable().c_str(), + previousGranule.end.printable().c_str()); + // must have retried and successfully nuked everything + return Void(); + } + ASSERT(totalState.size() >= 2); + + int total = totalState.size(); + int totalStarted = 0; + int totalDone = 0; + BlobGranuleSplitState currentState = BlobGranuleSplitState::Unknown; + for (auto& it : totalState) { + Tuple key = Tuple::unpack(it.key.removePrefix(blobGranuleSplitKeys.begin)); + ASSERT(key.getString(0) == previousGranule.begin); + ASSERT(key.getString(1) == previousGranule.end); + + BlobGranuleSplitState st = decodeBlobGranuleSplitValue(it.value); + ASSERT(st != BlobGranuleSplitState::Unknown); + if (st == BlobGranuleSplitState::Started) { + totalStarted++; + } else if (st == BlobGranuleSplitState::Done) { + totalDone++; + } + if (key.getString(2) == currentGranule.begin) { + ASSERT(currentState == BlobGranuleSplitState::Unknown); + currentState = st; + } + } + + ASSERT(currentState != BlobGranuleSplitState::Unknown); + + if (currentState < newState) { + printf("Updating granule [%s - %s) split state from [%s - %s) %d -> %d\n", + currentGranule.begin.printable().c_str(), + currentGranule.end.printable().c_str(), + previousGranule.begin.printable().c_str(), + previousGranule.end.printable().c_str(), + currentState, + newState); + + Tuple myStateTuple; + myStateTuple.append(previousGranule.begin).append(previousGranule.end).append(currentGranule.begin); + Key myStateKey = myStateTuple.getDataAsStandalone().withPrefix(blobGranuleSplitKeys.begin); + if (newState == BlobGranuleSplitState::Done && currentState == BlobGranuleSplitState::Assigned && + totalDone == total - 1) { + // we are the last one to change from Assigned -> Done, so everything can be cleaned up for the old change + // feed and splitting state + printf("[%s - %s) destroying old change feed %s and granule lock + split state for [%s - %s)\n", + currentGranule.begin.printable().c_str(), + currentGranule.end.printable().c_str(), + prevChangeFeedId.toString().c_str(), + previousGranule.begin.printable().c_str(), + previousGranule.end.printable().c_str()); + Key oldGranuleLockKey = granuleLockKey(previousGranule); + tr->destroyRangeFeed(KeyRef(prevChangeFeedId.toString())); + tr->clear(singleKeyRange(oldGranuleLockKey)); + tr->clear(currentRange); + } else { + if (newState == BlobGranuleSplitState::Assigned && currentState == BlobGranuleSplitState::Started && + totalStarted == 1) { + printf("[%s - %s) WOULD BE stopping old change feed %s for [%s - %s)\n", + currentGranule.begin.printable().c_str(), + currentGranule.end.printable().c_str(), + prevChangeFeedId.toString().c_str(), + previousGranule.begin.printable().c_str(), + previousGranule.end.printable().c_str()); + // FIXME: enable once implemented + // tr.stopChangeFeed(KeyRef(prevChangeFeedId.toString())); + } + tr->set(myStateKey, blobGranuleSplitValueFor(newState)); + } + } else { + printf("Ignoring granule [%s - %s) split state from [%s - %s) %d -> %d\n", + currentGranule.begin.printable().c_str(), + currentGranule.end.printable().c_str(), + previousGranule.begin.printable().c_str(), + previousGranule.end.printable().c_str(), + currentState, + newState); + } + + return Void(); +} + static Value getFileValue(std::string fname, int64_t offset, int64_t length) { Tuple fileValue; fileValue.append(fname).append(offset).append(length); @@ -175,7 +426,9 @@ ACTOR Future writeDeltaFile(BlobWorkerData* bwData, int64_t epoch, int64_t seqno, GranuleDeltas const* deltasToWrite, - Version currentDeltaVersion) { + Version currentDeltaVersion, + Optional oldChangeFeedDataComplete, + Optional oldChangeFeedId) { // TODO some sort of directory structure would be useful? state std::string fname = deterministicRandom()->randomUniqueID().toString() + "_T" + @@ -202,14 +455,27 @@ ACTOR Future writeDeltaFile(BlobWorkerData* bwData, wait(readAndCheckGranuleLock(tr, keyRange, epoch, seqno)); Tuple deltaFileKey; deltaFileKey.append(keyRange.begin).append(keyRange.end); - deltaFileKey.append(LiteralStringRef("delta")).append(currentDeltaVersion); + deltaFileKey.append(LiteralStringRef("D")).append(currentDeltaVersion); - tr->set(deltaFileKey.getDataAsStandalone().withPrefix(blobGranuleFileKeys.begin), - getFileValue(fname, 0, serialized.size())); + Key dfKey = deltaFileKey.getDataAsStandalone().withPrefix(blobGranuleFileKeys.begin); + tr->set(dfKey, getFileValue(fname, 0, serialized.size())); + + // FIXME: if previous granule present and delta file version >= previous change feed version, update the + // state here + if (oldChangeFeedDataComplete.present()) { + ASSERT(oldChangeFeedId.present()); + wait(updateGranuleSplitState(&tr->getTransaction(), + oldChangeFeedDataComplete.get(), + keyRange, + oldChangeFeedId.get(), + BlobGranuleSplitState::Done)); + } wait(tr->commit()); if (BW_DEBUG) { - printf("blob worker updated fdb with delta file %s of size %d at version %lld\n", + printf("Granule [%s - %s) updated fdb with delta file %s of size %d at version %lld\n", + keyRange.begin.printable().c_str(), + keyRange.end.printable().c_str(), fname.c_str(), serialized.size(), currentDeltaVersion); @@ -299,7 +565,7 @@ ACTOR Future writeSnapshot(BlobWorkerData* bwData, // TODO add conflict range for writes? state Tuple snapshotFileKey; snapshotFileKey.append(keyRange.begin).append(keyRange.end); - snapshotFileKey.append(LiteralStringRef("snapshot")).append(version); + snapshotFileKey.append(LiteralStringRef("S")).append(version); state Reference tr = makeReference(bwData->db); @@ -357,7 +623,7 @@ ACTOR Future dumpInitialSnapshotFromFDB(BlobWorkerData* bwData, R state Version readVersion = wait(tr->getReadVersion()); state PromiseStream rowsStream; state Future snapshotWriter = writeSnapshot( - bwData, metadata->keyRange, metadata->lockEpoch, metadata->lockSeqno, readVersion, rowsStream); + bwData, metadata->keyRange, metadata->originalEpoch, metadata->originalSeqno, readVersion, rowsStream); loop { // TODO: use streaming range read @@ -386,56 +652,62 @@ ACTOR Future dumpInitialSnapshotFromFDB(BlobWorkerData* bwData, R } } -ACTOR Future compactFromBlob(BlobWorkerData* bwData, Reference metadata) { +// files might not be the current set of files in metadata, in the case of doing the initial snapshot of a granule. +ACTOR Future compactFromBlob(BlobWorkerData* bwData, + Reference metadata, + GranuleFiles files) { if (BW_DEBUG) { printf("Compacting snapshot from blob for [%s - %s)\n", metadata->keyRange.begin.printable().c_str(), metadata->keyRange.end.printable().c_str()); } - ASSERT(!metadata->snapshotFiles.empty()); - ASSERT(!metadata->deltaFiles.empty()); - ASSERT(metadata->currentDeltas.empty()); - state Version version = metadata->deltaFiles.back().version; + ASSERT(!files.snapshotFiles.empty()); + ASSERT(!files.deltaFiles.empty()); + state Version version = files.deltaFiles.back().version; state Arena filenameArena; state BlobGranuleChunkRef chunk; - state Version snapshotVersion = metadata->snapshotFiles.back().version; - BlobFileIndex snapshotF = metadata->snapshotFiles.back(); + state int64_t compactBytesRead = 0; + state Version snapshotVersion = files.snapshotFiles.back().version; + BlobFileIndex snapshotF = files.snapshotFiles.back(); chunk.snapshotFile = BlobFilenameRef(filenameArena, snapshotF.filename, snapshotF.offset, snapshotF.length); - int deltaIdx = metadata->deltaFiles.size() - 1; - while (deltaIdx >= 0 && metadata->deltaFiles[deltaIdx].version > snapshotVersion) { + compactBytesRead += snapshotF.length; + int deltaIdx = files.deltaFiles.size() - 1; + while (deltaIdx >= 0 && files.deltaFiles[deltaIdx].version > snapshotVersion) { deltaIdx--; } deltaIdx++; - while (deltaIdx < metadata->deltaFiles.size()) { - BlobFileIndex deltaF = metadata->deltaFiles[deltaIdx]; + while (deltaIdx < files.deltaFiles.size()) { + BlobFileIndex deltaF = files.deltaFiles[deltaIdx]; chunk.deltaFiles.emplace_back_deep(filenameArena, deltaF.filename, deltaF.offset, deltaF.length); + compactBytesRead += deltaF.length; deltaIdx++; } chunk.includedVersion = version; if (BW_DEBUG) { - printf("Re-snapshotting [%s - %s) @ %lld\n", + printf("Re-snapshotting [%s - %s) @ %lld from blob\n", metadata->keyRange.begin.printable().c_str(), metadata->keyRange.end.printable().c_str(), version); - printf(" SnapshotFile:\n %s\n", chunk.snapshotFile.get().toString().c_str()); + /*printf(" SnapshotFile:\n %s\n", chunk.snapshotFile.get().toString().c_str()); printf(" DeltaFiles:\n"); for (auto& df : chunk.deltaFiles) { - printf(" %s\n", df.toString().c_str()); - } + printf(" %s\n", df.toString().c_str()); + }*/ } loop { try { state PromiseStream rowsStream; state Future snapshotWriter = writeSnapshot( - bwData, metadata->keyRange, metadata->lockEpoch, metadata->lockSeqno, version, rowsStream); - RangeResult newGranule = wait(readBlobGranule(chunk, metadata->keyRange, version, bwData->bstore, &bwData->stats)); - bwData->stats.bytesReadFromS3ForCompaction += newGranule.expectedSize(); + bwData, metadata->keyRange, metadata->originalEpoch, metadata->originalSeqno, version, rowsStream); + RangeResult newGranule = + wait(readBlobGranule(chunk, metadata->keyRange, version, bwData->bstore, &bwData->stats)); + bwData->stats.bytesReadFromS3ForCompaction += compactBytesRead; rowsStream.send(std::move(newGranule)); rowsStream.sendError(end_of_stream()); @@ -455,58 +727,274 @@ ACTOR Future compactFromBlob(BlobWorkerData* bwData, Reference> createRangeFeed(BlobWorkerData* bwData, KeyRange keyRange) { - state Key rangeFeedID = StringRef(deterministicRandom()->randomUniqueID().toString()); - state Transaction tr(bwData->db); - - loop { - try { - tr.setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS); - wait(tr.registerRangeFeed(rangeFeedID, keyRange)); - wait(tr.commit()); - return std::pair(rangeFeedID, tr.getCommittedVersion()); - } catch (Error& e) { - wait(tr.onError(e)); +// When reading from a prior change feed, the prior change feed may contain mutations that don't belong in the new +// granule. And, we only want to read the prior change feed up to the start of the new change feed. +static bool filterOldMutations(const KeyRange& range, + const Standalone>* oldMutations, + Standalone>* mutations, + Version maxVersion) { + Standalone> filteredMutations; + mutations->arena().dependsOn(range.arena()); + mutations->arena().dependsOn(oldMutations->arena()); + for (auto& delta : *oldMutations) { + if (delta.version >= maxVersion) { + return true; } + MutationsAndVersionRef filteredDelta; + filteredDelta.version = delta.version; + for (auto& m : delta.mutations) { + ASSERT(m.type == MutationRef::SetValue || m.type == MutationRef::ClearRange); + if (m.type == MutationRef::SetValue) { + if (m.param1 >= range.begin && m.param1 < range.end) { + filteredDelta.mutations.push_back(mutations->arena(), m); + } + } else { + if (m.param2 >= range.begin && m.param1 < range.end) { + // clamp clear range down to sub-range + MutationRef m2 = m; + if (range.begin > m.param1) { + m2.param1 = range.begin; + } + if (range.end < m.param2) { + m2.param2 = range.end; + } + filteredDelta.mutations.push_back(mutations->arena(), m2); + } + } + } + mutations->push_back(mutations->arena(), filteredDelta); } + return false; } // TODO a hack, eventually just have open end interval in range feed request? // Or maybe we want to cycle and start a new range feed stream every X million versions? -static const Version maxVersion = std::numeric_limits::max(); // updater for a single granule +// TODO: this is getting kind of large. Should try to split out this actor if it continues to grow? +// FIXME: handle errors here (forward errors) ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, Reference metadata) { - state PromiseStream>> rangeFeedStream; - state Future rangeFeedFuture; - try { - // before starting, make sure worker persists range assignment and acquires the granule lock - wait(metadata->assignFuture); - // TODO refactor creating range feed into the same transaction as above? + state PromiseStream>> oldChangeFeedStream; + state PromiseStream>> changeFeedStream; + state Future oldChangeFeedFuture; + state Future changeFeedFuture; + state GranuleChangeFeedInfo changeFeedInfo; + state bool readOldChangeFeed; + state bool lastFromOldChangeFeed = false; + state Optional oldChangeFeedDataComplete; + state Key cfKey; + state Optional oldCFKey; - // create range feed first so the version the SS start recording mutations <= the snapshot version - state std::pair rangeFeedData = wait(createRangeFeed(bwData, metadata->keyRange)); - if (BW_DEBUG) { - printf("Successfully created range feed %s for [%s - %s) @ %lld\n", - rangeFeedData.first.printable().c_str(), - metadata->keyRange.begin.printable().c_str(), - metadata->keyRange.end.printable().c_str(), - rangeFeedData.second); + try { + // set resume snapshot so it's not valid until we pause to ask the blob manager for a re-snapshot + metadata->resumeSnapshot.send(Void()); + + // before starting, make sure worker persists range assignment and acquires the granule lock + GranuleChangeFeedInfo _info = wait(metadata->assignFuture); + changeFeedInfo = _info; + cfKey = StringRef(changeFeedInfo.changeFeedId.toString()); + if (changeFeedInfo.prevChangeFeedId.present()) { + oldCFKey = StringRef(changeFeedInfo.prevChangeFeedId.get().toString()); } - BlobFileIndex newSnapshotFile = wait(dumpInitialSnapshotFromFDB(bwData, metadata)); - ASSERT(rangeFeedData.second <= newSnapshotFile.version); - metadata->snapshotFiles.push_back(newSnapshotFile); - metadata->lastWriteVersion = newSnapshotFile.version; - metadata->currentDeltaVersion = metadata->lastWriteVersion; - rangeFeedFuture = bwData->db->getRangeFeedStream( - rangeFeedStream, rangeFeedData.first, newSnapshotFile.version + 1, maxVersion, metadata->keyRange); + printf("Granule File Updater Starting for [%s - %s):\n", + metadata->keyRange.begin.printable().c_str(), + metadata->keyRange.end.printable().c_str()); + printf(" CFID: %s\n", changeFeedInfo.changeFeedId.toString().c_str()); + printf(" CF Start Version: %lld\n", changeFeedInfo.changeFeedStartVersion); + printf(" Previous Durable Version: %lld\n", changeFeedInfo.previousDurableVersion); + printf(" doSnapshot=%s\n", changeFeedInfo.doSnapshot ? "T" : "F"); + printf(" Prev CFID: %s\n", + changeFeedInfo.prevChangeFeedId.present() ? changeFeedInfo.prevChangeFeedId.get().toString().c_str() + : ""); + printf(" granuleSplitFrom=%s\n", changeFeedInfo.granuleSplitFrom.present() ? "T" : "F"); + printf(" blobFilesToSnapshot=%s\n", changeFeedInfo.blobFilesToSnapshot.present() ? "T" : "F"); + + // FIXME: handle reassigns by not doing a snapshot!! + state Version startVersion; + state BlobFileIndex newSnapshotFile; + + // FIXME: not true for reassigns + ASSERT(changeFeedInfo.doSnapshot); + if (!changeFeedInfo.doSnapshot) { + startVersion = changeFeedInfo.previousDurableVersion; + // TODO metadata.files = + } else { + if (changeFeedInfo.blobFilesToSnapshot.present()) { + BlobFileIndex fromBlob = + wait(compactFromBlob(bwData, metadata, changeFeedInfo.blobFilesToSnapshot.get())); + newSnapshotFile = fromBlob; + } else { + ASSERT(changeFeedInfo.previousDurableVersion == invalidVersion); + BlobFileIndex fromFDB = wait(dumpInitialSnapshotFromFDB(bwData, metadata)); + newSnapshotFile = fromFDB; + ASSERT(changeFeedInfo.changeFeedStartVersion <= fromFDB.version); + } + + startVersion = newSnapshotFile.version; + metadata->files.snapshotFiles.push_back(newSnapshotFile); + metadata->lastWriteVersion = newSnapshotFile.version; + metadata->currentDeltaVersion = metadata->lastWriteVersion; + } + + if (changeFeedInfo.prevChangeFeedId.present()) { + // FIXME: once we have empty versions, only include up to changeFeedInfo.changeFeedStartVersion in the read + // stream. Then we can just stop the old stream when we get end_of_stream from this and not handle the + // mutation version truncation stuff + ASSERT(changeFeedInfo.granuleSplitFrom.present()); + readOldChangeFeed = true; + oldChangeFeedFuture = bwData->db->getRangeFeedStream( + oldChangeFeedStream, oldCFKey.get(), startVersion + 1, MAX_VERSION, metadata->keyRange); + } else { + readOldChangeFeed = false; + changeFeedFuture = bwData->db->getRangeFeedStream( + changeFeedStream, cfKey, startVersion + 1, MAX_VERSION, metadata->keyRange); + } loop { + // TODO: handle empty versions here // TODO: Buggify delay in change feed stream - state Standalone> mutations = waitNext(rangeFeedStream.getFuture()); - for (auto& deltas : mutations) { + state Standalone> mutations; + if (readOldChangeFeed) { + Standalone> oldMutations = waitNext(oldChangeFeedStream.getFuture()); + if (filterOldMutations( + metadata->keyRange, &oldMutations, &mutations, changeFeedInfo.changeFeedStartVersion)) { + // if old change feed has caught up with where new one would start, finish last one and start new + // one + readOldChangeFeed = false; + Key cfKey = StringRef(changeFeedInfo.changeFeedId.toString()); + changeFeedFuture = bwData->db->getRangeFeedStream(changeFeedStream, + cfKey, + changeFeedInfo.changeFeedStartVersion, + MAX_VERSION, + metadata->keyRange); + oldChangeFeedFuture.cancel(); + lastFromOldChangeFeed = true; + + // now that old change feed is cancelled, clear out any mutations still in buffer by replacing + // promise stream + oldChangeFeedStream = PromiseStream>>(); + } + } else { + Standalone> newMutations = waitNext(changeFeedStream.getFuture()); + mutations = newMutations; + } + + // process mutations + for (MutationsAndVersionRef d : mutations) { + state MutationsAndVersionRef deltas = d; + ASSERT(deltas.version >= metadata->currentDeltaVersion); + // Write a new delta file IF we have enough bytes, and we have all of the previous version's stuff + // there to ensure no versions span multiple delta files. Check this by ensuring the version of this new + // delta is larger than the previous largest seen version + if (metadata->currentDeltaBytes >= SERVER_KNOBS->BG_DELTA_FILE_TARGET_BYTES && + deltas.version > metadata->currentDeltaVersion) { + if (BW_DEBUG) { + printf("Granule [%s - %s) flushing delta file after %d bytes\n", + metadata->keyRange.begin.printable().c_str(), + metadata->keyRange.end.printable().c_str(), + metadata->currentDeltaBytes); + } + TraceEvent("BlobGranuleDeltaFile", bwData->id) + .detail("GranuleStart", metadata->keyRange.begin) + .detail("GranuleEnd", metadata->keyRange.end) + .detail("Version", metadata->currentDeltaVersion); + BlobFileIndex newDeltaFile = wait(writeDeltaFile(bwData, + metadata->keyRange, + metadata->originalEpoch, + metadata->originalSeqno, + &metadata->currentDeltas, + metadata->currentDeltaVersion, + oldChangeFeedDataComplete, + changeFeedInfo.prevChangeFeedId)); + + // add new delta file + metadata->files.deltaFiles.push_back(newDeltaFile); + metadata->lastWriteVersion = metadata->currentDeltaVersion; + metadata->bytesInNewDeltaFiles += metadata->currentDeltaBytes; + + bwData->stats.mutationBytesBuffered -= metadata->currentDeltaBytes; + + // reset current deltas + metadata->deltaArena = Arena(); + metadata->currentDeltas = GranuleDeltas(); + metadata->currentDeltaBytes = 0; + + if (!readOldChangeFeed && !lastFromOldChangeFeed) { + if (BW_DEBUG) { + printf("Popping range feed %s at %lld\n\n", + changeFeedInfo.changeFeedId.toString().c_str(), + metadata->lastWriteVersion); + } + wait(bwData->db->popRangeFeedMutations(cfKey, metadata->lastWriteVersion)); + } + oldChangeFeedDataComplete.reset(); + + // if we just wrote a delta file, check if we need to compact here. + // exhaust old change feed before compacting - otherwise we could end up with an endlessly growing + // list + // of previous change feeds in the worst case. + if (metadata->bytesInNewDeltaFiles >= SERVER_KNOBS->BG_DELTA_BYTES_BEFORE_COMPACT && + !readOldChangeFeed && !lastFromOldChangeFeed) { + if (BW_DEBUG) { + printf("Granule [%s - %s) checking with BM for re-snapshot after %d bytes\n", + metadata->keyRange.begin.printable().c_str(), + metadata->keyRange.end.printable().c_str(), + metadata->bytesInNewDeltaFiles); + } + + TraceEvent("BlobGranuleSnapshotCheck", bwData->id) + .detail("GranuleStart", metadata->keyRange.begin) + .detail("GranuleEnd", metadata->keyRange.end) + .detail("Version", metadata->currentDeltaVersion); + + // Save these from the start so repeated requests are idempotent + // Need to retry in case response is dropped or manager changes. Eventually, a manager will + // either reassign the range with continue=true, or will revoke the range. But, we will keep the + // range open at this version for reads until that assignment change happens + metadata->resumeSnapshot.reset(); + state int64_t statusEpoch = metadata->continueEpoch; + state int64_t statusSeqno = metadata->continueSeqno; + loop { + bwData->currentManagerStatusStream.send( + GranuleStatusReply(metadata->keyRange, true, statusEpoch, statusSeqno)); + + Optional result = wait(timeout(metadata->resumeSnapshot.getFuture(), 1.0)); + if (result.present()) { + break; + } + if (BW_DEBUG) { + printf("Granule [%s - %s)\n, hasn't heard back from BM, re-sending status\n", + metadata->keyRange.begin.printable().c_str(), + metadata->keyRange.end.printable().c_str()); + } + } + + if (BW_DEBUG) { + printf("Granule [%s - %s) re-snapshotting after %d bytes\n", + metadata->keyRange.begin.printable().c_str(), + metadata->keyRange.end.printable().c_str(), + metadata->bytesInNewDeltaFiles); + } + TraceEvent("BlobGranuleSnapshotFile", bwData->id) + .detail("GranuleStart", metadata->keyRange.begin) + .detail("GranuleEnd", metadata->keyRange.end) + .detail("Version", metadata->currentDeltaVersion); + // TODO: this could read from FDB instead if it knew there was a large range clear at the end or + // it knew the granule was small, or something + BlobFileIndex newSnapshotFile = wait(compactFromBlob(bwData, metadata, metadata->files)); + + // add new snapshot file + metadata->files.snapshotFiles.push_back(newSnapshotFile); + metadata->lastWriteVersion = newSnapshotFile.version; + + // reset metadata + metadata->bytesInNewDeltaFiles = 0; + } + } + + // finally, after we optionally write delta and snapshot files, add new mutations to buffer if (!deltas.mutations.empty()) { metadata->currentDeltas.push_back_deep(metadata->deltaArena, deltas); for (auto& delta : deltas.mutations) { @@ -520,70 +1008,25 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, ReferencecurrentDeltaVersion <= deltas.version); metadata->currentDeltaVersion = deltas.version; - - // TODO handle version batch barriers - if (metadata->currentDeltaBytes >= SERVER_KNOBS->BG_DELTA_FILE_TARGET_BYTES && - metadata->currentDeltaVersion > metadata->lastWriteVersion) { - if (BW_DEBUG) { - printf("Granule [%s - %s) flushing delta file after %d bytes\n", - metadata->keyRange.begin.printable().c_str(), - metadata->keyRange.end.printable().c_str(), - metadata->currentDeltaBytes); - } - BlobFileIndex newDeltaFile = wait(writeDeltaFile(bwData, - metadata->keyRange, - metadata->lockEpoch, - metadata->lockSeqno, - &metadata->currentDeltas, - metadata->currentDeltaVersion)); - - // add new delta file - metadata->deltaFiles.push_back(newDeltaFile); - metadata->lastWriteVersion = metadata->currentDeltaVersion; - metadata->bytesInNewDeltaFiles += metadata->currentDeltaBytes; - - bwData->stats.mutationBytesBuffered -= metadata->currentDeltaBytes; - - // reset current deltas - metadata->deltaArena = Arena(); - metadata->currentDeltas = GranuleDeltas(); - metadata->currentDeltaBytes = 0; - - if (BW_DEBUG) { - printf("Popping range feed %s at %lld\n\n", - rangeFeedData.first.printable().c_str(), - metadata->lastWriteVersion); - } - wait(bwData->db->popRangeFeedMutations(rangeFeedData.first, metadata->lastWriteVersion)); - } - - if (metadata->bytesInNewDeltaFiles >= SERVER_KNOBS->BG_DELTA_BYTES_BEFORE_COMPACT) { - if (BW_DEBUG) { - printf("Granule [%s - %s) re-snapshotting after %d bytes\n", - metadata->keyRange.begin.printable().c_str(), - metadata->keyRange.end.printable().c_str(), - metadata->bytesInNewDeltaFiles); - } - // FIXME: instead of just doing new snapshot, it should offer shard back to blob manager and get - // reassigned - // TODO: this could read from FDB read previous snapshot + delta files instead if it knew there was - // a large range clear at the end or it knew the granule was small, or something - BlobFileIndex newSnapshotFile = wait(compactFromBlob(bwData, metadata)); - - // add new snapshot file - metadata->snapshotFiles.push_back(newSnapshotFile); - metadata->lastWriteVersion = newSnapshotFile.version; - - // reset metadata - metadata->bytesInNewDeltaFiles = 0; - } + } + if (lastFromOldChangeFeed) { + lastFromOldChangeFeed = false; + // set this so next delta file write updates granule split metadata to done + ASSERT(changeFeedInfo.granuleSplitFrom.present()); + oldChangeFeedDataComplete = changeFeedInfo.granuleSplitFrom; } } } catch (Error& e) { - printf("Granule file updater for [%s - %s) got error %s, exiting\n", - metadata->keyRange.begin.printable().c_str(), - metadata->keyRange.end.printable().c_str(), - e.name()); + if (BW_DEBUG) { + printf("Granule file updater for [%s - %s) got error %s, exiting\n", + metadata->keyRange.begin.printable().c_str(), + metadata->keyRange.end.printable().c_str(), + e.name()); + } + TraceEvent(SevError, "GranuleFileUpdaterError", bwData->id) + .detail("GranuleStart", metadata->keyRange.begin) + .detail("GranuleEnd", metadata->keyRange.end) + .error(e); // TODO in this case, need to update range mapping that it doesn't have the range, and/or try to re-"open" the // range if someone else doesn't have it throw e; @@ -645,8 +1088,9 @@ static void handleBlobGranuleFileRequest(BlobWorkerData* bwData, const BlobGranu StringRef(rep.arena, r.value().activeMetadata->keyRange.end)); // handle snapshot files - int i = metadata->snapshotFiles.size() - 1; - while (i >= 0 && metadata->snapshotFiles[i].version > req.readVersion) { + // TODO refactor the "find snapshot file" logic to GranuleFiles + int i = metadata->files.snapshotFiles.size() - 1; + while (i >= 0 && metadata->files.snapshotFiles[i].version > req.readVersion) { i--; } // if version is older than oldest snapshot file (or no snapshot files), throw too old @@ -656,33 +1100,33 @@ static void handleBlobGranuleFileRequest(BlobWorkerData* bwData, const BlobGranu printf("Oldest snapshot file for [%s - %s) is @ %lld, later than request version %lld\n", req.keyRange.begin.printable().c_str(), req.keyRange.end.printable().c_str(), - metadata->snapshotFiles.size() == 0 ? 0 : metadata->snapshotFiles[0].version, + metadata->files.snapshotFiles.size() == 0 ? 0 : metadata->files.snapshotFiles[0].version, req.readVersion); } req.reply.sendError(transaction_too_old()); return; } - BlobFileIndex snapshotF = metadata->snapshotFiles[i]; + BlobFileIndex snapshotF = metadata->files.snapshotFiles[i]; chunk.snapshotFile = BlobFilenameRef(rep.arena, snapshotF.filename, snapshotF.offset, snapshotF.length); - Version snapshotVersion = metadata->snapshotFiles[i].version; + Version snapshotVersion = metadata->files.snapshotFiles[i].version; // handle delta files - i = metadata->deltaFiles.size() - 1; + i = metadata->files.deltaFiles.size() - 1; // skip delta files that are too new - while (i >= 0 && metadata->deltaFiles[i].version > req.readVersion) { + while (i >= 0 && metadata->files.deltaFiles[i].version > req.readVersion) { i--; } - if (i < metadata->deltaFiles.size() - 1) { + if (i < metadata->files.deltaFiles.size() - 1) { i++; } // only include delta files after the snapshot file int j = i; - while (j >= 0 && metadata->deltaFiles[j].version > snapshotVersion) { + while (j >= 0 && metadata->files.deltaFiles[j].version > snapshotVersion) { j--; } j++; while (j <= i) { - BlobFileIndex deltaF = metadata->deltaFiles[j]; + BlobFileIndex deltaF = metadata->files.deltaFiles[j]; chunk.deltaFiles.emplace_back_deep(rep.arena, deltaF.filename, deltaF.offset, deltaF.length); bwData->stats.readReqDeltaBytesReturned += deltaF.length; j++; @@ -690,7 +1134,7 @@ static void handleBlobGranuleFileRequest(BlobWorkerData* bwData, const BlobGranu // new deltas (if version is larger than version of last delta file) // FIXME: do trivial key bounds here if key range is not fully contained in request key range - if (!metadata->deltaFiles.size() || req.readVersion >= metadata->deltaFiles.back().version) { + if (!metadata->files.deltaFiles.size() || req.readVersion >= metadata->files.deltaFiles.back().version) { rep.arena.dependsOn(metadata->deltaArena); for (auto& delta : metadata->currentDeltas) { if (delta.version <= req.readVersion) { @@ -708,45 +1152,107 @@ static void handleBlobGranuleFileRequest(BlobWorkerData* bwData, const BlobGranu req.reply.send(rep); } -// TODO list of key ranges in the future to batch -ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t epoch, int64_t seqno) { - state Reference tr = makeReference(bwData->db); - state Key lockKey = granuleLockKey(keyRange); +// FIXME: in split, need to persist version of created change feed so if worker immediately fails afterwards, new worker +// picking up the splitting shard knows where the change feed handoff point is. OR need to have change feed return +// end_of_stream when it knows it has nothing up to the specified end version, and use the commit takeover version as +// the end version. If it sealed successfully there would trivially be nothing between the seal version and the new +// commit takeover version. You'd need to start the new change feed at the seal version though, not the commit takeover +// version. +ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwData, AssignBlobRangeRequest req) { + ASSERT(!req.continueAssignment); + state Transaction tr(bwData->db); + state Key lockKey = granuleLockKey(req.keyRange); + state GranuleChangeFeedInfo info; + info.changeFeedId = deterministicRandom()->randomUniqueID(); + if (BW_DEBUG) { + printf("%s persisting assignment [%s - %s)\n", + bwData->id.toString().c_str(), + req.keyRange.begin.printable().c_str(), + req.keyRange.end.printable().c_str()); + } loop { try { - tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS); - tr->setOption(FDBTransactionOptions::PRIORITY_SYSTEM_IMMEDIATE); + tr.setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS); + tr.setOption(FDBTransactionOptions::PRIORITY_SYSTEM_IMMEDIATE); - Optional prevLockValue = wait(tr->get(lockKey)); + // FIXME: could add list of futures and do the different parts that are disjoint in parallel? + info.changeFeedStartVersion = invalidVersion; + Optional prevLockValue = wait(tr.get(lockKey)); if (prevLockValue.present()) { - std::pair prevOwner = decodeBlobGranuleLockValue(prevLockValue.get()); - acquireGranuleLock(epoch, seqno, prevOwner); - } // else we are first, no need to check for owner conflict + std::tuple prevOwner = decodeBlobGranuleLockValue(prevLockValue.get()); + acquireGranuleLock(req.managerEpoch, req.managerSeqno, std::get<0>(prevOwner), std::get<1>(prevOwner)); + info.changeFeedId = std::get<2>(prevOwner); + info.doSnapshot = false; - tr->set(lockKey, blobGranuleLockValueFor(epoch, seqno)); + ASSERT(info.changeFeedId == UID()); - wait(krmSetRangeCoalescing( - tr, blobGranuleMappingKeys.begin, keyRange, KeyRange(allKeys), blobGranuleMappingValueFor(bwData->id))); + /*info.existingFiles = wait(loadPreviousFiles(&tr, req.keyRange)); + info.previousDurableVersion = info.existingFiles.get().deltaFiles.empty() + ? info.existingFiles.get().snapshotFiles.back().version + : info.existingFiles.get().deltaFiles.back().version;*/ + // FIXME: Handle granule reassignments! + ASSERT(false); - wait(tr->commit()); - - if (BW_DEBUG) { - printf("Blob worker %s persisted key range [%s - %s)\n", - bwData->id.toString().c_str(), - keyRange.begin.printable().c_str(), - keyRange.end.printable().c_str()); + } else { + // else we are first, no need to check for owner conflict + wait(tr.registerRangeFeed(StringRef(info.changeFeedId.toString()), req.keyRange)); + info.doSnapshot = true; + info.previousDurableVersion = invalidVersion; } - return Void(); + + tr.set(lockKey, blobGranuleLockValueFor(req.managerEpoch, req.managerSeqno, info.changeFeedId)); + wait(krmSetRange(&tr, blobGranuleMappingKeys.begin, req.keyRange, blobGranuleMappingValueFor(bwData->id))); + + // If anything in previousGranules, need to do the handoff logic and set ret.previousChangeFeedId, and the + // previous durable version will come from the previous granules + if (!req.previousGranules.empty()) { + // TODO change this for merge + ASSERT(req.previousGranules.size() == 1); + Optional prevGranuleLockValue = wait(tr.get(granuleLockKey(req.previousGranules[0]))); + + ASSERT(prevGranuleLockValue.present()); + + std::tuple prevGranuleLock = + decodeBlobGranuleLockValue(prevGranuleLockValue.get()); + info.prevChangeFeedId = std::get<2>(prevGranuleLock); + + wait(updateGranuleSplitState(&tr, + req.previousGranules[0], + req.keyRange, + info.prevChangeFeedId.get(), + BlobGranuleSplitState::Assigned)); + + // FIXME: store this somewhere useful for time travel reads + GranuleFiles prevFiles = wait(loadPreviousFiles(&tr, req.previousGranules[0])); + ASSERT(!prevFiles.snapshotFiles.empty() || !prevFiles.deltaFiles.empty()); + info.granuleSplitFrom = req.previousGranules[0]; + info.blobFilesToSnapshot = prevFiles; + info.previousDurableVersion = info.blobFilesToSnapshot.get().deltaFiles.empty() + ? info.blobFilesToSnapshot.get().snapshotFiles.back().version + : info.blobFilesToSnapshot.get().deltaFiles.back().version; + + // FIXME: need to handle takeover of a splitting range! If snapshot and/or deltas found for new range, + // don't snapshot + } + // else: FIXME: If nothing in previousGranules, previous durable version is max of previous snapshot version + // and previous delta version. If neither present, need to do a snapshot at the start. + // Assumes for now that this isn't a takeover, so nothing to do here + wait(tr.commit()); + + TraceEvent("BlobWorkerPersistedAssignment", bwData->id) + .detail("GranuleStart", req.keyRange.begin) + .detail("GranuleEnd", req.keyRange.end); + + if (info.changeFeedStartVersion == invalidVersion) { + info.changeFeedStartVersion = tr.getCommittedVersion(); + } + return info; } catch (Error& e) { - if (BW_DEBUG) { - printf("Persisting key range [%s - %s) for blob worker %s got error %s\n", - keyRange.begin.printable().c_str(), - keyRange.end.printable().c_str(), - bwData->id.toString().c_str(), - e.name()); + if (e.code() == error_code_granule_assignment_conflict) { + throw e; } - wait(tr->onError(e)); + wait(tr.onError(e)); } } } @@ -755,18 +1261,13 @@ static GranuleRangeMetadata constructActiveBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t epoch, int64_t seqno) { - if (BW_DEBUG) { - printf("Creating new worker metadata for range [%s - %s)\n", - keyRange.begin.printable().c_str(), - keyRange.end.printable().c_str()); - } Reference newMetadata = makeReference(); newMetadata->keyRange = keyRange; - newMetadata->lockEpoch = epoch; - newMetadata->lockSeqno = seqno; - newMetadata->assignFuture = persistAssignWorkerRange(bwData, keyRange, epoch, seqno); - newMetadata->fileUpdaterFuture = blobGranuleUpdateFiles(bwData, newMetadata); + newMetadata->originalEpoch = epoch; + newMetadata->originalSeqno = seqno; + newMetadata->continueEpoch = epoch; + newMetadata->continueSeqno = seqno; return GranuleRangeMetadata(epoch, seqno, newMetadata); } @@ -793,7 +1294,16 @@ static bool newerRangeAssignment(GranuleRangeMetadata oldMetadata, int64_t epoch // in, it is a no-op, but updates the sequence number. Similarly, if a worker gets an assign message for any range that // already has a higher sequence number, that range was either revoked, or revoked and then re-assigned. Either way, // this assignment is no longer valid. -static void changeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t epoch, int64_t seqno, bool active) { + +// Returns future to wait on to ensure prior work of other granules is done before responding to the manager with a +// successful assignment And if the change produced a new granule that needs to start doing work, returns the new +// granule so that the caller can start() it with the appropriate starting state. +static std::pair, Reference> changeBlobRange(BlobWorkerData* bwData, + KeyRange keyRange, + int64_t epoch, + int64_t seqno, + bool active, + bool disposeOnCleanup) { if (BW_DEBUG) { printf("Changing range for [%s - %s): %s @ (%lld, %lld)\n", keyRange.begin.printable().c_str(), @@ -810,6 +1320,8 @@ static void changeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t e // Insert the current range. // Re-insert all newer ranges over the current range. + std::vector> futures; + std::vector> newerRanges; auto ranges = bwData->granuleMetadata.intersectingRanges(keyRange); @@ -818,7 +1330,10 @@ static void changeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t e // applied the same assignment twice, make idempotent ASSERT(r.begin() == keyRange.begin); ASSERT(r.end() == keyRange.end); - return; + if (r.value().activeMetadata.isValid()) { + futures.push_back(success(r.value().activeMetadata->assignFuture)); + } + return std::pair(waitForAll(futures), Reference()); // already applied, nothing to do } bool thisAssignmentNewer = newerRangeAssignment(r.value(), epoch, seqno); if (r.value().activeMetadata.isValid() && thisAssignmentNewer) { @@ -830,7 +1345,7 @@ static void changeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t e r.value().lastEpoch, r.value().lastSeqno); } - r.value().activeMetadata->cancel(); + futures.push_back(r.value().activeMetadata->cancel(disposeOnCleanup)); r.value().activeMetadata.clear(); } else if (!thisAssignmentNewer) { // this assignment is outdated, re-insert it over the current range @@ -863,14 +1378,43 @@ static void changeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t e } bwData->granuleMetadata.insert(it.first, it.second); } + + return std::pair(waitForAll(futures), newMetadata.activeMetadata); } -static void handleAssignedRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t epoch, int64_t seqno) { - changeBlobRange(bwData, keyRange, epoch, seqno, true); -} +static bool resumeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t epoch, int64_t seqno) { + auto existingRange = bwData->granuleMetadata.rangeContaining(keyRange.begin); + // if range boundaries don't match, or this (epoch, seqno) is old or the granule is inactive, ignore + if (keyRange.begin != existingRange.begin() || keyRange.end != existingRange.end() || + existingRange.value().lastEpoch > epoch || + (existingRange.value().lastEpoch == epoch && existingRange.value().lastSeqno > seqno) || + !existingRange.value().activeMetadata.isValid()) { -static void handleRevokedRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t epoch, int64_t seqno) { - changeBlobRange(bwData, keyRange, epoch, seqno, false); + printf("BW %s got out of date continue range for [%s - %s) @ (%lld, %lld). Currently [%s - %s) @ (%lld, " + "%lld): %s\n", + bwData->id.toString().c_str(), + existingRange.begin().printable().c_str(), + existingRange.end().printable().c_str(), + existingRange.value().lastEpoch, + existingRange.value().lastSeqno, + keyRange.begin.printable().c_str(), + keyRange.end.printable().c_str(), + epoch, + seqno, + existingRange.value().activeMetadata.isValid() ? "T" : "F"); + + return false; + } + if (existingRange.value().lastEpoch != epoch || existingRange.value().lastSeqno != seqno) { + // update the granule metadata map, and the continueEpoch/seqno. Saves an extra transaction + existingRange.value().lastEpoch = epoch; + existingRange.value().lastSeqno = seqno; + existingRange.value().activeMetadata->continueEpoch = epoch; + existingRange.value().activeMetadata->continueSeqno = seqno; + existingRange.value().activeMetadata->resume(); + } + // else we already processed this continue, do nothing + return true; } ACTOR Future registerBlobWorker(BlobWorkerData* bwData, BlobWorkerInterface interf) { @@ -898,6 +1442,62 @@ ACTOR Future registerBlobWorker(BlobWorkerData* bwData, BlobWorkerInterfac } } +// TODO might want to separate this out for valid values for range assignments vs read requests +namespace { +bool canReplyWith(Error e) { + switch (e.code()) { + case error_code_transaction_too_old: + case error_code_future_version: // not thrown yet + case error_code_wrong_shard_server: + case error_code_process_behind: // not thrown yet + // TODO should we reply with granule_assignment_conflict? + return true; + default: + return false; + }; +} +} // namespace + +ACTOR Future handleRangeAssign(BlobWorkerData* bwData, AssignBlobRangeRequest req) { + try { + if (req.continueAssignment) { + resumeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno); + } else { + // FIXME: wait to reply unless worker confirms it should own range and takes out lock? + state std::pair, Reference> futureAndNewGranule = + changeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno, true, false); + + wait(futureAndNewGranule.first); + + if (futureAndNewGranule.second.isValid()) { + wait(futureAndNewGranule.second->start(bwData, req)); + } + } + req.reply.send(AssignBlobRangeReply(true)); + return Void(); + } catch (Error& e) { + printf("AssignRange got error %s\n", e.name()); + if (canReplyWith(e)) { + req.reply.sendError(e); + } + throw; + } +} + +ACTOR Future handleRangeRevoke(BlobWorkerData* bwData, RevokeBlobRangeRequest req) { + try { + wait(changeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno, false, req.dispose).first); + req.reply.send(AssignBlobRangeReply(true)); + return Void(); + } catch (Error& e) { + printf("RevokeRange got error %s\n", e.name()); + if (canReplyWith(e)) { + req.reply.sendError(e); + } + throw; + } +} + // TODO need to version assigned ranges with read version of txn the range was read and use that in // handleAssigned/Revoked from to prevent out-of-order assignments/revokes from the blob manager from getting ranges in // an incorrect state @@ -949,40 +1549,62 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, Reference self.currentManagerEpoch) { - self.currentManagerEpoch = req.managerEpoch; - if (BW_DEBUG) { - printf("BW %s found new manager epoch %lld\n", - self.id.toString().c_str(), - self.currentManagerEpoch); - } - } + assignReq.reply.send(AssignBlobRangeReply(false)); + } + } + when(RevokeBlobRangeRequest _req = waitNext(bwInterf.revokeBlobRangeRequest.getFuture())) { + state RevokeBlobRangeRequest revokeReq = _req; + --self.stats.numRangesAssigned; + if (BW_DEBUG) { + printf("Worker %s revoked range [%s - %s) @ (%lld, %lld):\n dispose=%s\n", + self.id.toString().c_str(), + revokeReq.keyRange.begin.printable().c_str(), + revokeReq.keyRange.end.printable().c_str(), + revokeReq.managerEpoch, + revokeReq.managerSeqno, + revokeReq.dispose ? "T" : "F"); + } - // TODO with range versioning, need to persist only after it's confirmed - changeBlobRange(&self, req.keyRange, req.managerEpoch, req.managerSeqno, req.isAssign); - req.isAssign ? ++self.stats.numRangesAssigned : --self.stats.numRangesAssigned; - req.reply.send(AssignBlobRangeReply(true)); + if (self.managerEpochOk(revokeReq.managerEpoch)) { + addActor.send(handleRangeRevoke(&self, revokeReq)); + } else { + revokeReq.reply.send(AssignBlobRangeReply(false)); } } when(wait(collection)) { diff --git a/fdbserver/storageserver.actor.cpp b/fdbserver/storageserver.actor.cpp index 967f894a82..2449b552d5 100644 --- a/fdbserver/storageserver.actor.cpp +++ b/fdbserver/storageserver.actor.cpp @@ -1561,6 +1561,16 @@ ACTOR Future getRangeFeedMutations(StorageServer* data, RangeFee if (data->version.get() < req.begin) { wait(data->version.whenAtLeast(req.begin)); } + // TODO REMOVE this super hacky fix once evan finishes change feeds, the above wait doesn't work apparently + state int waitCnt = 0; + while (data->uidRangeFeed.count(req.rangeID) == 0) { + wait(delay(0.05)); + waitCnt++; + if ((waitCnt & (waitCnt - 1)) == 0) { + printf("Waiting for change feed %s %d times\n", req.rangeID.printable().c_str(), waitCnt); + } + } + ASSERT(data->uidRangeFeed.count(req.rangeID) > 0); auto& feedInfo = data->uidRangeFeed[req.rangeID]; /*printf("SS processing range feed req %s for version [%lld - %lld)\n", req.rangeID.printable().c_str(), @@ -3878,6 +3888,12 @@ private: rangeFeedInfo->range = rangeFeedRange; rangeFeedInfo->id = rangeFeedId; rangeFeedInfo->emptyVersion = currentVersion - 1; + printf("SS %s creating change feed %s for [%s - %s) at version %lld\n", + data->thisServerID.toString().c_str(), + rangeFeedId.toString().c_str(), + rangeFeedRange.begin.printable().c_str(), + rangeFeedRange.end.printable().c_str(), + currentVersion); data->uidRangeFeed[rangeFeedId] = rangeFeedInfo; auto rs = data->keyRangeFeed.modify(rangeFeedRange); for (auto r = rs.begin(); r != rs.end(); ++r) { diff --git a/fdbserver/workloads/BlobGranuleVerifier.actor.cpp b/fdbserver/workloads/BlobGranuleVerifier.actor.cpp index a8f430266f..69e75e5a0e 100644 --- a/fdbserver/workloads/BlobGranuleVerifier.actor.cpp +++ b/fdbserver/workloads/BlobGranuleVerifier.actor.cpp @@ -59,10 +59,11 @@ struct BlobGranuleVerifierWorkload : TestWorkload { BlobGranuleVerifierWorkload(WorkloadContext const& wcx) : TestWorkload(wcx) { doSetup = !clientId; // only do this on the "first" client + // FIXME: don't do the delay in setup, as that delays the start of all workloads minDelay = getOption(options, LiteralStringRef("minDelay"), 0.0); - maxDelay = getOption(options, LiteralStringRef("minDelay"), 60.0); + maxDelay = getOption(options, LiteralStringRef("maxDelay"), 0.0); testDuration = getOption(options, LiteralStringRef("testDuration"), 120.0); - timeTravelLimit = getOption(options, LiteralStringRef("timeTravelLimit"), 60.0); + timeTravelLimit = getOption(options, LiteralStringRef("timeTravelLimit"), testDuration); timeTravelBufferSize = getOption(options, LiteralStringRef("timeTravelBufferSize"), 100000000); threads = getOption(options, LiteralStringRef("threads"), 1); ASSERT(threads >= 1); @@ -96,6 +97,7 @@ struct BlobGranuleVerifierWorkload : TestWorkload { wait(krmSetRange(tr, blobRangeKeys.begin, KeyRange(normalKeys), LiteralStringRef("1"))); wait(tr->commit()); printf("Successfully set up blob granule range for normalKeys\n"); + TraceEvent("BlobGranuleVerifierSetup"); return Void(); } catch (Error& e) { wait(tr->onError(e)); @@ -106,7 +108,6 @@ struct BlobGranuleVerifierWorkload : TestWorkload { std::string description() const override { return "BlobGranuleVerifier"; } Future setup(Database const& cx) override { if (doSetup) { - /// TODO make only one client do this!!! others wait double initialDelay = deterministicRandom()->random01() * (maxDelay - minDelay) + minDelay; printf("BGW setup initial delay of %.3f\n", initialDelay); return setUpBlobRange(cx, delay(initialDelay)); @@ -223,6 +224,7 @@ struct BlobGranuleVerifierWorkload : TestWorkload { state std::map timeTravelChecks; state int64_t timeTravelChecksMemory = 0; + TraceEvent("BlobGranuleVerifierStart"); printf("BGV thread starting\n"); // wait for first set of ranges to be loaded @@ -264,16 +266,18 @@ struct BlobGranuleVerifierWorkload : TestWorkload { self->bytesRead += fdb.first.expectedSize(); self->initialReads++; - // TODO increase frequency a lot!! just for initial testing - wait(poisson(&last, 5.0)); - // wait(poisson(&last, 0.1)); } catch (Error& e) { - printf("BGVerifier got error %s\n", e.name()); if (e.code() == error_code_operation_cancelled) { - return Void(); + throw; + } + if (e.code() != error_code_transaction_too_old && e.code() != error_code_wrong_shard_server) { + printf("BGVerifier got unexpected error %s\n", e.name()); } self->errors++; } + // TODO increase frequency a lot!! just for initial testing + wait(poisson(&last, 5.0)); + // wait(poisson(&last, 0.1)); } } @@ -329,6 +333,7 @@ struct BlobGranuleVerifierWorkload : TestWorkload { printf(" %lld rows\n", self->rowsRead); printf(" %lld bytes\n", self->bytesRead); printf(" %d final granule checks\n", checks); + TraceEvent("BlobGranuleVerifierChecked"); return self->mismatches == 0 && checks > 0; } diff --git a/tests/slow/BlobGranuleCorrectnessLarge.toml b/tests/slow/BlobGranuleCorrectnessLarge.toml new file mode 100644 index 0000000000..499deb4d2c --- /dev/null +++ b/tests/slow/BlobGranuleCorrectnessLarge.toml @@ -0,0 +1,21 @@ +[[test]] +testTitle = 'BlobGranuleCorrectnessTestLarge' + + [[test.workload]] + testName = 'ReadWrite' + testDuration = 200.0 + transactionsPerSecond = 1000 + writesPerTransactionA = 0 + readsPerTransactionA = 10 + writesPerTransactionB = 10 + readsPerTransactionB = 1 + alpha = 0.5 + nodeCount = 2000000 + valueBytes = 128 + discardEdgeMeasurements = false + warmingDelay = 10.0 + setup = false + + [[test.workload]] + testName = 'BlobGranuleVerifier' + testDuration = 200.0