diff --git a/fdbclient/BlobWorkerInterface.h b/fdbclient/BlobWorkerInterface.h index eb2be3623e..eafa8eb338 100644 --- a/fdbclient/BlobWorkerInterface.h +++ b/fdbclient/BlobWorkerInterface.h @@ -35,6 +35,8 @@ struct BlobWorkerInterface { RequestStream assignBlobRangeRequest; RequestStream revokeBlobRangeRequest; RequestStream granuleStatusStreamRequest; + RequestStream haltBlobWorker; + struct LocalityData locality; UID myId; @@ -57,6 +59,7 @@ struct BlobWorkerInterface { assignBlobRangeRequest, revokeBlobRangeRequest, granuleStatusStreamRequest, + haltBlobWorker, locality, myId); } @@ -182,4 +185,20 @@ struct GranuleStatusStreamRequest { } }; +struct HaltBlobWorkerRequest { + constexpr static FileIdentifier file_identifier = 1985879; + UID requesterID; + ReplyPromise reply; + + int64_t managerEpoch; + + HaltBlobWorkerRequest() {} + explicit HaltBlobWorkerRequest(int64_t managerEpoch, UID uid) : requesterID(uid), managerEpoch(managerEpoch) {} + + template + void serialize(Ar& ar) { + serializer(ar, managerEpoch, requesterID, reply); + } +}; + #endif diff --git a/fdbrpc/fdbrpc.h b/fdbrpc/fdbrpc.h index c2641b95b2..b4a3606d94 100644 --- a/fdbrpc/fdbrpc.h +++ b/fdbrpc/fdbrpc.h @@ -481,12 +481,14 @@ public: const Endpoint& getEndpoint() const { return queue->getEndpoint(TaskPriority::ReadSocket); } bool operator==(const ReplyPromiseStream& rhs) const { return queue == rhs.queue; } + bool operator!=(const ReplyPromiseStream& rhs) const { return !(*this == rhs); } + bool isEmpty() const { return !queue->isReady(); } uint32_t size() const { return queue->size(); } // Must be called on the server before sending results on the stream to ratelimit the amount of data outstanding to // the client - Future onReady() { + Future onReady() const { ASSERT(queue->acknowledgements.bytesLimit > 0); if (queue->acknowledgements.failures.isError()) { return queue->acknowledgements.failures.getError(); diff --git a/fdbserver/BlobManager.actor.cpp b/fdbserver/BlobManager.actor.cpp index f1ad4ffc3f..a524a6a4e2 100644 --- a/fdbserver/BlobManager.actor.cpp +++ b/fdbserver/BlobManager.actor.cpp @@ -298,7 +298,7 @@ ACTOR Future doRangeAssignment(BlobManagerData* bmData, RangeAssignment as if (BM_DEBUG) { printf("BM %s %s range [%s - %s) @ (%lld, %lld)\n", - workerID.toString().c_str(), + bmData->id.toString().c_str(), assignment.isAssign ? "assigning" : "revoking", assignment.keyRange.begin.printable().c_str(), assignment.keyRange.end.printable().c_str(), @@ -318,6 +318,11 @@ ACTOR Future doRangeAssignment(BlobManagerData* bmData, RangeAssignment as req.managerEpoch = bmData->epoch; req.managerSeqno = seqNo; req.continueAssignment = assignment.assign.get().continueAssignment; + + // if that worker isn't alive anymore, add the range back into the stream + if (bmData->workersById.count(workerID) == 0) { + throw no_more_servers(); + } AssignBlobRangeReply _rep = wait(bmData->workersById[workerID].assignBlobRangeRequest.getReply(req)); rep = _rep; } else { @@ -331,8 +336,13 @@ ACTOR Future doRangeAssignment(BlobManagerData* bmData, RangeAssignment as req.managerSeqno = seqNo; req.dispose = assignment.revoke.get().dispose; - AssignBlobRangeReply _rep = wait(bmData->workersById[workerID].revokeBlobRangeRequest.getReply(req)); - rep = _rep; + // if that worker isn't alive anymore, this is a noop + if (bmData->workersById.count(workerID)) { + AssignBlobRangeReply _rep = wait(bmData->workersById[workerID].revokeBlobRangeRequest.getReply(req)); + rep = _rep; + } else { + return Void(); + } } if (!rep.epochOk) { if (BM_DEBUG) { @@ -349,7 +359,8 @@ ACTOR Future doRangeAssignment(BlobManagerData* bmData, RangeAssignment as if (BM_DEBUG) { printf("BM got error assigning range [%s - %s) to worker %s, requeueing\n", assignment.keyRange.begin.printable().c_str(), - assignment.keyRange.end.printable().c_str()); + assignment.keyRange.end.printable().c_str(), + workerID.toString().c_str()); } // re-send revoke to queue to handle range being un-assigned from that worker before the new one RangeAssignment revokeOld; @@ -417,7 +428,11 @@ ACTOR Future rangeAssigner(BlobManagerData* bmData) { workerId = assignment.worker.present() ? assignment.worker.get() : pickWorkerForAssign(bmData); bmData->workerAssignments.insert(assignment.keyRange, workerId); - bmData->workerStats[workerId].numGranulesAssigned += 1; + + ASSERT(bmData->workerStats.count(workerId)); + if (!assignment.assign.get().continueAssignment) { + bmData->workerStats[workerId].numGranulesAssigned += 1; + } // FIXME: if range is assign, have some sort of semaphore for outstanding assignments so we don't assign // a ton ranges at once and blow up FDB with reading initial snapshots. @@ -432,9 +447,13 @@ 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; - if (!assignment.worker.present() || assignment.worker.get() == it.value()) - bmData->addActor.send(doRangeAssignment(bmData, assignment, it.value(), seqNo)); + + if (bmData->workerStats.count(it.value())) { + bmData->workerStats[it.value()].numGranulesAssigned -= 1; + } + + // revoke the range for the worker that owns it, not the worker specified in the revoke + bmData->addActor.send(doRangeAssignment(bmData, assignment, it.value(), seqNo)); } bmData->workerAssignments.insert(assignment.keyRange, UID()); @@ -701,6 +720,49 @@ ACTOR Future maybeSplitRange(BlobManagerData* bmData, UID currentWorkerId, return Void(); } +void killBlobWorker(BlobManagerData* bmData, BlobWorkerInterface bwInterf) { + UID bwId = bwInterf.id(); + + // Remove blob worker from stats map so that when we try to find a worker to takeover the range, + // the one we just killed isn't considered. + // Remove it from workersById also since otherwise that worker addr will remain excluded + // when we try to recruit new blob workers. + bmData->workerStats.erase(bwId); + bmData->workersById.erase(bwId); + + // for every range owned by this blob worker, we want to + // - send a revoke request for that range + // - add the range back to the stream of ranges to be assigned + if (BM_DEBUG) { + printf("Taking back ranges from BW %s\n", bwId.toString().c_str()); + } + for (auto& it : bmData->workerAssignments.ranges()) { + if (it.cvalue() == bwId) { + // Send revoke request + RangeAssignment raRevoke; + raRevoke.isAssign = false; + raRevoke.keyRange = it.range(); + raRevoke.revoke = RangeRevokeData(false); + bmData->rangesToAssign.send(raRevoke); + + // Add range back into the stream of ranges to be assigned + RangeAssignment raAssign; + raAssign.isAssign = true; + raAssign.worker = Optional(); + raAssign.keyRange = it.range(); + raAssign.assign = RangeAssignmentData(false); // not a continue + bmData->rangesToAssign.send(raAssign); + } + } + + // Send halt to blob worker, with no expectation of hearing back + if (BM_DEBUG) { + printf("Sending halt to BW %s\n", bwId.toString().c_str()); + } + bmData->addActor.send( + brokenPromiseToNever(bwInterf.haltBlobWorker.getReply(HaltBlobWorkerRequest(bmData->epoch, bmData->id)))); +} + ACTOR Future monitorBlobWorkerStatus(BlobManagerData* bmData, BlobWorkerInterface bwInterf) { state KeyRangeMap> lastSeenSeqno; // outer loop handles reconstructing stream if it got a retryable error @@ -711,6 +773,7 @@ ACTOR Future monitorBlobWorkerStatus(BlobManagerData* bmData, BlobWorkerIn // read from stream until worker fails (should never get explicit end_of_stream) loop { GranuleStatusReply rep = waitNext(statusStream.getFuture()); + if (BM_DEBUG) { printf("BM %lld got status of [%s - %s) @ (%lld, %lld) from BW %s: %s\n", bmData->epoch, @@ -723,7 +786,8 @@ ACTOR Future monitorBlobWorkerStatus(BlobManagerData* bmData, BlobWorkerIn } if (rep.epoch > bmData->epoch) { if (BM_DEBUG) { - printf("BM heard from BW that there is a new manager with higher epoch\n"); + printf("BM heard from BW %s that there is a new manager with higher epoch\n", + bwInterf.id().toString().c_str()); } if (bmData->iAmReplaced.canBeSet()) { bmData->iAmReplaced.send(Void()); @@ -734,8 +798,13 @@ ACTOR Future monitorBlobWorkerStatus(BlobManagerData* bmData, BlobWorkerIn // 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 + // only evaluate for split if this worker currently owns the granule in this blob manager's mapping + auto currGranuleAssignment = bmData->workerAssignments.rangeContaining(rep.granuleRange.begin); + if (!(currGranuleAssignment.begin() == rep.granuleRange.begin && + currGranuleAssignment.end() == rep.granuleRange.end && + currGranuleAssignment.cvalue() == bwInterf.id())) { + continue; + } auto lastReqForGranule = lastSeenSeqno.rangeContaining(rep.granuleRange.begin); if (rep.granuleRange.begin == lastReqForGranule.begin() && @@ -792,16 +861,10 @@ ACTOR Future monitorBlobWorker(BlobManagerData* bmData, BlobWorkerInterfac 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()); } TraceEvent("BlobWorkerFailed", bmData->id).detail("BlobWorkerID", bwInterf.id()); - // get all of its ranges - // send revoke request to get back all its ranges - // send halt (look at rangeMover) - // send all its ranges to assignranges stream - return Void(); } when(wait(monitorStatus)) { ASSERT(false); @@ -820,6 +883,20 @@ ACTOR Future monitorBlobWorker(BlobManagerData* bmData, BlobWorkerInterfac TraceEvent(SevError, "BWMonitoringFailed", bmData->id).detail("BlobWorkerID", bwInterf.id()).error(e); throw e; } + + // kill the blob worker + killBlobWorker(bmData, bwInterf); + + // Trigger recruitment for a new blob worker + if (BM_DEBUG) { + printf("Restarting recruitment to replace dead BW %s\n", bwInterf.id().toString().c_str()); + } + bmData->restartRecruiting.trigger(); + + if (BM_DEBUG) { + printf("No longer monitoring BW %s\n", bwInterf.id().toString().c_str()); + } + return Void(); } // TODO this is only for chaos testing right now!! REMOVE LATER @@ -854,7 +931,6 @@ ACTOR Future rangeMover(BlobManagerData* bmData) { RangeAssignment revokeOld; revokeOld.isAssign = false; revokeOld.keyRange = randomRange.range(); - revokeOld.worker = randomRange.value(); revokeOld.revoke = RangeRevokeData(false); bmData->rangesToAssign.send(revokeOld); @@ -937,8 +1013,8 @@ ACTOR Future initializeBlobWorker(BlobManagerData* self, RecruitBlobWorker if (newBlobWorker.present()) { BlobWorkerInterface bwi = newBlobWorker.get().interf; - self->workersById.insert({ bwi.id(), bwi }); - self->workerStats.insert({ bwi.id(), BlobWorkerStats() }); + self->workersById[bwi.id()] = bwi; + self->workerStats[bwi.id()] = BlobWorkerStats(); self->addActor.send(monitorBlobWorker(self, bwi)); TraceEvent("BMRecruiting") diff --git a/fdbserver/BlobWorker.actor.cpp b/fdbserver/BlobWorker.actor.cpp index 0782cbddab..45605e8934 100644 --- a/fdbserver/BlobWorker.actor.cpp +++ b/fdbserver/BlobWorker.actor.cpp @@ -34,6 +34,7 @@ #include "fdbserver/MutationTracking.h" #include "fdbserver/WaitFailure.h" #include "flow/Arena.h" +#include "flow/Error.h" #include "flow/IRandom.h" #include "flow/actorcompiler.h" // has to be last include #include "flow/flow.h" @@ -41,8 +42,6 @@ #define BW_DEBUG true #define BW_REQUEST_DEBUG false -// FIXME: change all BlobWorkerData* to Reference to avoid segfaults if core loop gets error - // TODO add comments + documentation struct BlobFileIndex { Version version; @@ -73,28 +72,22 @@ struct GranuleChangeFeedInfo { Optional existingFiles; }; -// 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 { KeyRange keyRange; GranuleFiles files; - GranuleDeltas currentDeltas; + GranuleDeltas currentDeltas; // only contain deltas in pendingDeltaVersion + 1, bufferedDeltaVersion // TODO get rid of this and do Reference>? Arena deltaArena; uint64_t bytesInNewDeltaFiles = 0; uint64_t bufferedDeltaBytes = 0; - NotifiedVersion bufferedDeltaVersion; - Version pendingDeltaVersion = 0; - NotifiedVersion durableDeltaVersion; - NotifiedVersion durableSnapshotVersion; + // for client to know when it is safe to read a certain version and from where (check waitForVersion) + NotifiedVersion bufferedDeltaVersion; // largest delta version in currentDeltas (including empty versions) + Version pendingDeltaVersion = 0; // largest version in progress writing to s3/fdb + NotifiedVersion durableDeltaVersion; // largest version persisted in s3/fdb + NotifiedVersion durableSnapshotVersion; // same as delta vars, except for snapshots Version pendingSnapshotVersion = 0; AsyncVar rollbackCount; @@ -104,68 +97,50 @@ struct GranuleMetadata : NonCopyable, ReferenceCounted { int64_t continueEpoch; int64_t continueSeqno; - Future assignFuture; - Future fileUpdaterFuture; Promise resumeSnapshot; + + // used to coordinate granule file updater is done Promise cancelled; Promise readable; AssignBlobRangeRequest originalReq; - Future start(BlobWorkerData* bwData, AssignBlobRangeRequest req) { - originalReq = 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. - Future cancel(bool dispose) { - if (cancelled.canBeSet()) { - // Could have been cancelled already by rollback or error in BGUpdateFiles - cancelled.send(Void()); - } - assignFuture.cancel(); - fileUpdaterFuture.cancel(); - - if (dispose) { - // FIXME: implement dispose! - return delay(0.1); - } - return Future(Void()); - } }; -// for a range that may or may not be set +// TODO: rename this struct struct GranuleRangeMetadata { int64_t lastEpoch; int64_t lastSeqno; Reference activeMetadata; + Future assignFuture; + Future fileUpdaterFuture; + GranuleRangeMetadata() : lastEpoch(0), lastSeqno(0) {} GranuleRangeMetadata(int64_t epoch, int64_t seqno, Reference activeMetadata) : lastEpoch(epoch), lastSeqno(seqno), activeMetadata(activeMetadata) {} }; -struct BlobWorkerData { +// FIXME: there is a reference cycle here. BWData has GranuleRangeMetadata objects in a map, +// but each of those has a future to a forever-running actor which has a reference to BWData. +// To fix this, we should only pass the necessary, specfic fields of BWData to those actors +// rather than the reference to BWData itself. +struct BlobWorkerData : NonCopyable, ReferenceCounted { UID id; Database db; BlobWorkerStats stats; + PromiseStream> addActor; + LocalityData locality; int64_t currentManagerEpoch = -1; - ReplyPromiseStream currentManagerStatusStream; + AsyncVar> currentManagerStatusStream; // FIXME: refactor out the parts of this that are just for interacting with blob stores from the backup business // logic @@ -221,7 +196,8 @@ static void checkGranuleLock(int64_t epoch, int64_t seqno, int64_t ownerEpoch, i // sanity check - lock value should never go backwards because of acquireGranuleLock /* printf( - "Checking granule lock: \n mine: (%lld, %lld)\n owner: (%lld, %lld)\n", epoch, seqno, ownerEpoch, ownerSeqno); + "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)); @@ -322,21 +298,23 @@ ACTOR Future loadPreviousFiles(Transaction* tr, KeyRange keyRange) // 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 +// 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 +// 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 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. +// 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, @@ -411,8 +389,8 @@ ACTOR Future updateGranuleSplitState(Transaction* tr, 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 + // we are the last one to change from Assigned -> Done, so everything can be cleaned up for the old + // change feed and splitting state if (BW_DEBUG) { printf("[%s - %s) destroying old change feed %s and granule lock + split state for [%s - %s)\n", currentGranule.begin.printable().c_str(), @@ -477,11 +455,12 @@ static Value getFileValue(std::string fname, int64_t offset, int64_t length) { return fileValue.getDataAsStandalone(); } -// writeDelta file writes speculatively in the common case to optimize throughput. It creates the s3 object even though -// the data in it may not yet be committed, and even though previous delta fiels with lower versioned data may still be -// in flight. The synchronization happens after the s3 file is written, but before we update the FDB index of what files -// exist. Before updating FDB, we ensure the version is committed and all previous delta files have updated FDB. -ACTOR Future writeDeltaFile(BlobWorkerData* bwData, +// writeDelta file writes speculatively in the common case to optimize throughput. It creates the s3 object even +// though the data in it may not yet be committed, and even though previous delta fiels with lower versioned data +// may still be in flight. The synchronization happens after the s3 file is written, but before we update the FDB +// index of what files exist. Before updating FDB, we ensure the version is committed and all previous delta files +// have updated FDB. +ACTOR Future writeDeltaFile(Reference bwData, KeyRange keyRange, int64_t epoch, int64_t seqno, @@ -517,6 +496,7 @@ ACTOR Future writeDeltaFile(BlobWorkerData* bwData, wait(objectFile->append(serialized.begin(), serialized.size())); wait(objectFile->finish()); + state int numIterations = 0; try { // before updating FDB, wait for the delta file version to be committed and previous delta files to finish while (bwData->knownCommittedVersion.get() < currentDeltaVersion) { @@ -543,8 +523,8 @@ ACTOR Future writeDeltaFile(BlobWorkerData* bwData, 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 + // 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(), @@ -570,6 +550,7 @@ ACTOR Future writeDeltaFile(BlobWorkerData* bwData, } return BlobFileIndex(currentDeltaVersion, fname, 0, serialized.size()); } catch (Error& e) { + numIterations++; wait(tr->onError(e)); } } @@ -578,6 +559,13 @@ ACTOR Future writeDeltaFile(BlobWorkerData* bwData, throw e; } + // if commit failed the first time due to granule assignment conflict (which is non-retryable), + // then the file key was persisted and we should delete it. Otherwise, the commit failed + // for some other reason and the key wasn't persisted, so we should just propogate the error + if (numIterations != 1 || e.code() != error_code_granule_assignment_conflict) { + throw e; + } + // FIXME: only delete if key doesn't exist if (BW_DEBUG) { printf("deleting s3 delta file %s after error %s\n", fname.c_str(), e.name()); @@ -589,7 +577,7 @@ ACTOR Future writeDeltaFile(BlobWorkerData* bwData, } } -ACTOR Future writeSnapshot(BlobWorkerData* bwData, +ACTOR Future writeSnapshot(Reference bwData, KeyRange keyRange, int64_t epoch, int64_t seqno, @@ -663,6 +651,7 @@ ACTOR Future writeSnapshot(BlobWorkerData* bwData, snapshotFileKey.append(LiteralStringRef("S")).append(version); state Reference tr = makeReference(bwData->db); + state int numIterations = 0; try { loop { @@ -674,6 +663,7 @@ ACTOR Future writeSnapshot(BlobWorkerData* bwData, wait(tr->commit()); break; } catch (Error& e) { + numIterations++; wait(tr->onError(e)); } } @@ -682,6 +672,13 @@ ACTOR Future writeSnapshot(BlobWorkerData* bwData, throw e; } + // if commit failed the first time due to granule assignment conflict (which is non-retryable), + // then the file key was persisted and we should delete it. Otherwise, the commit failed + // for some other reason and the key wasn't persisted, so we should just propogate the error + if (numIterations != 1 || e.code() != error_code_granule_assignment_conflict) { + throw e; + } + // FIXME: only delete if key doesn't exist if (BW_DEBUG) { printf("deleting s3 snapshot file %s after error %s\n", fname.c_str(), e.name()); @@ -707,7 +704,8 @@ ACTOR Future writeSnapshot(BlobWorkerData* bwData, return BlobFileIndex(version, fname, 0, serialized.size()); } -ACTOR Future dumpInitialSnapshotFromFDB(BlobWorkerData* bwData, Reference metadata) { +ACTOR Future dumpInitialSnapshotFromFDB(Reference bwData, + Reference metadata) { if (BW_DEBUG) { printf("Dumping snapshot from FDB for [%s - %s)\n", metadata->keyRange.begin.printable().c_str(), @@ -755,7 +753,7 @@ ACTOR Future dumpInitialSnapshotFromFDB(BlobWorkerData* bwData, R } // 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, +ACTOR Future compactFromBlob(Reference bwData, Reference metadata, GranuleFiles files) { wait(delay(0, TaskPriority::BlobWorkerUpdateStorage)); @@ -821,8 +819,8 @@ ACTOR Future compactFromBlob(BlobWorkerData* bwData, DEBUG_KEY_RANGE("BlobWorkerBlobSnapshot", version, metadata->keyRange, bwData->id); return f; } catch (Error& e) { - // TODO better error handling eventually - should retry unless the error is because another worker took over - // the range + // TODO better error handling eventually - should retry unless the error is because another worker took + // over the range if (BW_DEBUG) { printf("Compacting snapshot from blob for [%s - %s) got error %s\n", metadata->keyRange.begin.printable().c_str(), @@ -874,7 +872,7 @@ static bool filterOldMutations(const KeyRange& range, return false; } -ACTOR Future handleCompletedDeltaFile(BlobWorkerData* bwData, +ACTOR Future handleCompletedDeltaFile(Reference bwData, Reference metadata, BlobFileIndex completedDeltaFile, Key cfKey, @@ -998,8 +996,8 @@ static Version doGranuleRollback(Reference metadata, metadata->currentDeltas.resize(metadata->deltaArena, mIdx); - // delete all deltas in rollback range, but we can optimize here to just skip the uncommitted mutations directly - // and immediately pop the rollback out of inProgress + // delete all deltas in rollback range, but we can optimize here to just skip the uncommitted mutations + // directly and immediately pop the rollback out of inProgress metadata->bufferedDeltaVersion.set(rollbackVersion); cfRollbackVersion = mutationVersion; } @@ -1026,7 +1024,9 @@ static Version doGranuleRollback(Reference metadata, // 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) { +ACTOR Future blobGranuleUpdateFiles(Reference bwData, + Reference metadata, + Future assignFuture) { state PromiseStream>> oldChangeFeedStream; state PromiseStream>> changeFeedStream; state Future inFlightBlobSnapshot; @@ -1051,7 +1051,7 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, ReferenceresumeSnapshot.send(Void()); // before starting, make sure worker persists range assignment and acquires the granule lock - GranuleChangeFeedInfo _info = wait(metadata->assignFuture); + GranuleChangeFeedInfo _info = wait(assignFuture); changeFeedInfo = _info; wait(delay(0, TaskPriority::BlobWorkerUpdateStorage)); @@ -1131,9 +1131,9 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, Referencereadable.send(Void()); 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 + // 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 // FIXME: filtering on key range != change feed range doesn't work ASSERT(changeFeedInfo.granuleSplitFrom.present()); @@ -1197,8 +1197,8 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, ReferencekeyRange, &oldMutations, &mutations, changeFeedInfo.changeFeedStartVersion)) { - // if old change feed has caught up with where new one would start, finish last one and start new - // one + // if old change feed has caught up with where new one would start, finish last one and start + // new one Key cfKey = StringRef(changeFeedInfo.changeFeedId.toString()); changeFeedFuture = bwData->db->getChangeFeedStream(changeFeedStream, @@ -1223,8 +1223,8 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, Reference= metadata->bufferedDeltaVersion.get()); // 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 + // 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->bufferedDeltaBytes >= SERVER_KNOBS->BG_DELTA_FILE_TARGET_BYTES && deltas.version > metadata->bufferedDeltaVersion.get()) { if (BW_DEBUG) { @@ -1283,8 +1283,8 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, ReferencebufferedDeltaBytes = 0; // 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. + // exhaust old change feed before compacting - otherwise we could end up with an endlessly + // growing list of previous change feeds in the worst case. snapshotEligible = true; } @@ -1292,7 +1292,6 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, ReferencebytesInNewDeltaFiles >= SERVER_KNOBS->BG_DELTA_BYTES_BEFORE_COMPACT && !readOldChangeFeed) { - if (BW_DEBUG && (inFlightBlobSnapshot.isValid() || !inFlightDeltaFiles.empty())) { printf("Granule [%s - %s) ready to re-snapshot, waiting for outstanding %d snapshot and %d " "deltas to " @@ -1341,18 +1340,29 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, ReferencecontinueEpoch; 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; + loop { + try { + wait(bwData->currentManagerStatusStream.get().onReady()); + bwData->currentManagerStatusStream.get().send( + GranuleStatusReply(metadata->keyRange, true, statusEpoch, statusSeqno)); + break; + } catch (Error& e) { + printf("manager stream was changed\n"); + wait(bwData->currentManagerStatusStream.onChange()); + } } - // FIXME: re-trigger this loop if blob manager status stream changes + + choose { + when(wait(metadata->resumeSnapshot.getFuture())) { break; } + when(wait(delay(1.0))) {} + when(wait(bwData->currentManagerStatusStream.onChange())) {} + } + if (BW_DEBUG) { - printf("Granule [%s - %s)\n, hasn't heard back from BM, re-sending status\n", + printf("Granule [%s - %s)\n, hasn't heard back from BM in BW %s, re-sending status\n", metadata->keyRange.begin.printable().c_str(), - metadata->keyRange.end.printable().c_str()); + metadata->keyRange.end.printable().c_str(), + bwData->id.toString().c_str()); } } @@ -1379,8 +1389,8 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, ReferencebytesInNewDeltaFiles = 0; } else if (snapshotEligible && metadata->bytesInNewDeltaFiles >= SERVER_KNOBS->BG_DELTA_BYTES_BEFORE_COMPACT) { - // if we're in the old change feed case and can't snapshot but we have enough data to, don't queue - // too many delta files in parallel + // if we're in the old change feed case and can't snapshot but we have enough data to, don't + // queue too many delta files in parallel while (inFlightDeltaFiles.size() > 10) { if (BW_DEBUG) { printf("[%s - %s) Waiting on delta file b/c old change feed\n", @@ -1526,7 +1536,6 @@ ACTOR Future blobGranuleUpdateFiles(BlobWorkerData* bwData, Reference blobGranuleUpdateFiles(BlobWorkerData* bwData, Referenceid).detail("Granule", metadata->keyRange).error(e); if (granuleCanRetry(e)) { - // explicitly cancel all outstanding write futures BEFORE updating promise stream, to ensure they can't - // update files after the re-assigned granule acquires the lock + // explicitly cancel all outstanding write futures BEFORE updating promise stream, to ensure they + // can't update files after the re-assigned granule acquires the lock inFlightBlobSnapshot.cancel(); for (auto& f : inFlightDeltaFiles) { f.future.cancel(); @@ -1619,7 +1628,8 @@ static Future waitForVersion(Reference metadata, Version // if we don't have to wait for change feed version to catch up or wait for any pending file writes to complete, // nothing to do - /*printf(" [%s - %s) waiting for %lld\n readable:%s\n bufferedDelta=%lld\n pendingDelta=%lld\n " + /* + printf(" [%s - %s) waiting for %lld\n readable:%s\n bufferedDelta=%lld\n pendingDelta=%lld\n " "durableDelta=%lld\n pendingSnapshot=%lld\n durableSnapshot=%lld\n", metadata->keyRange.begin.printable().c_str(), metadata->keyRange.end.printable().c_str(), @@ -1629,7 +1639,8 @@ static Future waitForVersion(Reference metadata, Version metadata->pendingDeltaVersion, metadata->durableDeltaVersion.get(), metadata->pendingSnapshotVersion, - metadata->durableSnapshotVersion.get());*/ + metadata->durableSnapshotVersion.get()); + */ if (metadata->readable.isSet() && v <= metadata->bufferedDeltaVersion.get() && (v <= metadata->durableDeltaVersion.get() || @@ -1642,7 +1653,7 @@ static Future waitForVersion(Reference metadata, Version return waitForVersionActor(metadata, v); } -ACTOR Future handleBlobGranuleFileRequest(BlobWorkerData* bwData, BlobGranuleFileRequest req) { +ACTOR Future handleBlobGranuleFileRequest(Reference bwData, BlobGranuleFileRequest req) { try { // TODO REMOVE in api V2 ASSERT(req.beginVersion == 0); @@ -1664,6 +1675,7 @@ ACTOR Future handleBlobGranuleFileRequest(BlobWorkerData* bwData, BlobGran req.keyRange.begin.printable().c_str(), req.keyRange.end.printable().c_str()); } + throw wrong_shard_server(); } granules.push_back(r.value().activeMetadata); @@ -1821,7 +1833,8 @@ ACTOR Future handleBlobGranuleFileRequest(BlobWorkerData* bwData, BlobGran return Void(); } -ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwData, AssignBlobRangeRequest req) { +ACTOR Future persistAssignWorkerRange(Reference bwData, + AssignBlobRangeRequest req) { ASSERT(!req.continueAssignment); state Transaction tr(bwData->db); state Key lockKey = granuleLockKey(req.keyRange); @@ -1861,13 +1874,13 @@ ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwD info.previousDurableVersion = info.existingFiles.get().deltaFiles.back().version; } - // for the non-splitting cases, this doesn't need to be 100% accurate, it just needs to be smaller - // than the next delta file write. + // for the non-splitting cases, this doesn't need to be 100% accurate, it just needs to be + // smaller than the next delta file write. info.changeFeedStartVersion = info.previousDurableVersion; } else { // else we are first, no need to check for owner conflict - // FIXME: use actual 16 bytes of UID instead of converting it to 32 character string and then that - // to bytes + // FIXME: use actual 16 bytes of UID instead of converting it to 32 character string and then + // that to bytes info.changeFeedId = deterministicRandom()->randomUniqueID(); wait(tr.registerChangeFeed(StringRef(info.changeFeedId.toString()), req.keyRange)); info.doSnapshot = true; @@ -1882,8 +1895,8 @@ ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwD Optional parentGranulesValue = wait(tr.get(historyKey.getDataAsStandalone().withPrefix(blobGranuleHistoryKeys.begin))); - // 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 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 (parentGranulesValue.present()) { state Standalone> parentGranules = decodeBlobGranuleHistoryValue(parentGranulesValue.get()); @@ -1930,8 +1943,8 @@ ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwD req.keyRange, info.prevChangeFeedId.get(), BlobGranuleSplitState::Assigned)); - // change feed was created as part of this transaction, changeFeedStartVersion will be set - // later + // change feed was created as part of this transaction, changeFeedStartVersion will be + // set later } else { ASSERT(false); } @@ -1971,7 +1984,16 @@ ACTOR Future persistAssignWorkerRange(BlobWorkerData* bwD } } -static GranuleRangeMetadata constructActiveBlobRange(BlobWorkerData* bwData, +ACTOR Future start(Reference bwData, GranuleRangeMetadata* meta, AssignBlobRangeRequest req) { + ASSERT(meta->activeMetadata.isValid()); + meta->activeMetadata->originalReq = req; + meta->assignFuture = persistAssignWorkerRange(bwData, req); + meta->fileUpdaterFuture = blobGranuleUpdateFiles(bwData, meta->activeMetadata, meta->assignFuture); + wait(success(meta->assignFuture)); + return Void(); +} + +static GranuleRangeMetadata constructActiveBlobRange(Reference bwData, KeyRange keyRange, int64_t epoch, int64_t seqno) { @@ -2013,13 +2035,13 @@ static bool newerRangeAssignment(GranuleRangeMetadata oldMetadata, int64_t epoch // 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, - bool selfReassign) { +ACTOR Future changeBlobRange(Reference bwData, + KeyRange keyRange, + int64_t epoch, + int64_t seqno, + bool active, + bool disposeOnCleanup, + bool selfReassign) { if (BW_DEBUG) { printf("%s range for [%s - %s): %s @ (%lld, %lld)\n", selfReassign ? "Re-assigning" : "Changing", @@ -2036,12 +2058,21 @@ static std::pair, Reference> changeBlobRange(BlobW // older range, cancel it if it is active. Insert the current range. Re-insert all newer ranges over the current // range. - std::vector> futures; + state std::vector> futures; - std::vector> newerRanges; + state std::vector> newerRanges; auto ranges = bwData->granuleMetadata.intersectingRanges(keyRange); + bool alreadyAssigned = false; for (auto& r : ranges) { + if (!active) { + if (r.value().activeMetadata.isValid() && r.value().activeMetadata->cancelled.canBeSet()) { + if (BW_DEBUG) { + printf("Cancelling activeMetadata\n"); + } + r.value().activeMetadata->cancelled.send(Void()); + } + } bool thisAssignmentNewer = newerRangeAssignment(r.value(), epoch, seqno); if (r.value().lastEpoch == epoch && r.value().lastSeqno == seqno) { ASSERT(r.begin() == keyRange.begin); @@ -2050,12 +2081,13 @@ static std::pair, Reference> changeBlobRange(BlobW if (selfReassign) { thisAssignmentNewer = true; } else { + printf("same assignment\n"); // applied the same assignment twice, make idempotent if (r.value().activeMetadata.isValid()) { - futures.push_back(success(r.value().activeMetadata->assignFuture)); + futures.push_back(success(r.value().assignFuture)); } - return std::pair(waitForAll(futures), - Reference()); // already applied, nothing to do + alreadyAssigned = true; + break; } } @@ -2068,7 +2100,6 @@ static std::pair, Reference> changeBlobRange(BlobW r.value().lastEpoch, r.value().lastSeqno); } - 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 @@ -2076,10 +2107,16 @@ static std::pair, Reference> changeBlobRange(BlobW } } + if (alreadyAssigned) { + wait(waitForAll(futures)); // already applied, nothing to do + return false; + } + // if range is active, and isn't surpassed by a newer range already, insert an active range GranuleRangeMetadata newMetadata = (active && newerRanges.empty()) ? constructActiveBlobRange(bwData, keyRange, epoch, seqno) : constructInactiveBlobRange(epoch, seqno); + bwData->granuleMetadata.insert(keyRange, newMetadata); if (BW_DEBUG) { printf("Inserting new range [%s - %s): %s @ (%lld, %lld)\n", @@ -2102,10 +2139,12 @@ static std::pair, Reference> changeBlobRange(BlobW bwData->granuleMetadata.insert(it.first, it.second); } - return std::pair(waitForAll(futures), newMetadata.activeMetadata); + printf("returning from changeblobrange"); + wait(waitForAll(futures)); + return true; } -static bool resumeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t epoch, int64_t seqno) { +static bool resumeBlobRange(Reference 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() || @@ -2142,7 +2181,7 @@ static bool resumeBlobRange(BlobWorkerData* bwData, KeyRange keyRange, int64_t e return true; } -ACTOR Future registerBlobWorker(BlobWorkerData* bwData, BlobWorkerInterface interf) { +ACTOR Future registerBlobWorker(Reference bwData, BlobWorkerInterface interf) { state Reference tr = makeReference(bwData->db); loop { tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS); @@ -2167,18 +2206,20 @@ ACTOR Future registerBlobWorker(BlobWorkerData* bwData, BlobWorkerInterfac } } -ACTOR Future handleRangeAssign(BlobWorkerData* bwData, AssignBlobRangeRequest req, bool isSelfReassign) { +ACTOR Future handleRangeAssign(Reference bwData, + AssignBlobRangeRequest req, + bool isSelfReassign) { try { if (req.continueAssignment) { resumeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno); } else { - state std::pair, Reference> futureAndNewGranule = - changeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno, true, false, isSelfReassign); + bool shouldStart = wait( + changeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno, true, false, isSelfReassign)); - wait(futureAndNewGranule.first); - - if (futureAndNewGranule.second.isValid()) { - wait(futureAndNewGranule.second->start(bwData, req)); + if (shouldStart) { + auto m = bwData->granuleMetadata.rangeContaining(req.keyRange.begin); + ASSERT(m.begin() == req.keyRange.begin && m.end() == req.keyRange.end); + wait(start(bwData, &m.value(), req)); } } if (!isSelfReassign) { @@ -2199,14 +2240,15 @@ ACTOR Future handleRangeAssign(BlobWorkerData* bwData, AssignBlobRangeRequ req.reply.sendError(e); } } + throw; } } -ACTOR Future handleRangeRevoke(BlobWorkerData* bwData, RevokeBlobRangeRequest req) { +ACTOR Future handleRangeRevoke(Reference bwData, RevokeBlobRangeRequest req) { try { - wait( - changeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno, false, req.dispose, false).first); + bool _shouldStart = + wait(changeBlobRange(bwData, req.keyRange, req.managerEpoch, req.managerSeqno, false, req.dispose, false)); req.reply.send(AssignBlobRangeReply(true)); return Void(); } catch (Error& e) { @@ -2229,7 +2271,7 @@ ACTOR Future handleRangeRevoke(BlobWorkerData* bwData, RevokeBlobRangeRequ // uncommitted data. This means we must ensure the data is actually committed before "committing" those writes in // the blob granule. The simplest way to do this is to have the blob worker do a periodic GRV, which is guaranteed // to be an earlier committed version. -ACTOR Future runCommitVersionChecks(BlobWorkerData* bwData) { +ACTOR Future runCommitVersionChecks(Reference bwData) { state Transaction tr(bwData->db); loop { // only do grvs to get committed version if we need it to persist delta files @@ -2262,9 +2304,12 @@ ACTOR Future runCommitVersionChecks(BlobWorkerData* bwData) { ACTOR Future blobWorker(BlobWorkerInterface bwInterf, ReplyPromise recruitReply, Reference const> dbInfo) { - state BlobWorkerData self(bwInterf.id(), openDBOnServer(dbInfo, TaskPriority::DefaultEndpoint, LockAware::True)); - self.id = bwInterf.id(); - self.locality = bwInterf.locality; + state Reference self( + new BlobWorkerData(bwInterf.id(), openDBOnServer(dbInfo, TaskPriority::DefaultEndpoint, LockAware::True))); + self->id = bwInterf.id(); + self->locality = bwInterf.locality; + + state Future collection = actorCollection(self->addActor.getFuture()); if (BW_DEBUG) { printf("Initializing blob worker s3 stuff\n"); @@ -2275,19 +2320,19 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, if (BW_DEBUG) { printf("BW constructing simulated backup container\n"); } - self.bstore = BackupContainerFileSystem::openContainerFS("file://fdbblob/"); + self->bstore = BackupContainerFileSystem::openContainerFS("file://fdbblob/"); } else { if (BW_DEBUG) { printf("BW constructing backup container from %s\n", SERVER_KNOBS->BG_URL.c_str()); } - self.bstore = BackupContainerFileSystem::openContainerFS(SERVER_KNOBS->BG_URL); + self->bstore = BackupContainerFileSystem::openContainerFS(SERVER_KNOBS->BG_URL); if (BW_DEBUG) { printf("BW constructed backup container\n"); } } // register the blob worker to the system keyspace - wait(registerBlobWorker(&self, bwInterf)); + wait(registerBlobWorker(self, bwInterf)); } catch (Error& e) { if (BW_DEBUG) { printf("BW got backup container init error %s\n", e.name()); @@ -2307,11 +2352,8 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, rep.interf = bwInterf; recruitReply.send(rep); - state PromiseStream> addActor; - state Future collection = actorCollection(addActor.getFuture()); - - addActor.send(waitFailureServer(bwInterf.waitFailure.getFuture())); - addActor.send(runCommitVersionChecks(&self)); + self->addActor.send(waitFailureServer(bwInterf.waitFailure.getFuture())); + self->addActor.send(runCommitVersionChecks(self)); try { loop choose { @@ -2319,25 +2361,28 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, /*printf("Got blob granule request [%s - %s)\n", req.keyRange.begin.printable().c_str(), req.keyRange.end.printable().c_str());*/ - ++self.stats.readRequests; - ++self.stats.activeReadRequests; - addActor.send(handleBlobGranuleFileRequest(&self, req)); + ++self->stats.readRequests; + ++self->stats.activeReadRequests; + self->addActor.send(handleBlobGranuleFileRequest(self, req)); } - when(GranuleStatusStreamRequest req = waitNext(bwInterf.granuleStatusStreamRequest.getFuture())) { - if (self.managerEpochOk(req.managerEpoch)) { + when(state GranuleStatusStreamRequest req = waitNext(bwInterf.granuleStatusStreamRequest.getFuture())) { + if (self->managerEpochOk(req.managerEpoch)) { if (BW_DEBUG) { - printf("Worker %s got new granule status endpoint\n", self.id.toString().c_str()); + printf("Worker %s got new granule status endpoint\n", self->id.toString().c_str()); } - self.currentManagerStatusStream = req.reply; + // req.reply is marked const unless you mark req as `state`?!?!? + // TODO: pick a reasonable byte limit instead of just piggy-backing + req.reply.setByteLimit(SERVER_KNOBS->RANGESTREAM_LIMIT_BYTES); + self->currentManagerStatusStream.set(req.reply); } } when(AssignBlobRangeRequest _req = waitNext(bwInterf.assignBlobRangeRequest.getFuture())) { - ++self.stats.rangeAssignmentRequests; - --self.stats.numRangesAssigned; + ++self->stats.rangeAssignmentRequests; + --self->stats.numRangesAssigned; state AssignBlobRangeRequest assignReq = _req; if (BW_DEBUG) { printf("Worker %s assigned range [%s - %s) @ (%lld, %lld):\n continue=%s\n", - self.id.toString().c_str(), + self->id.toString().c_str(), assignReq.keyRange.begin.printable().c_str(), assignReq.keyRange.end.printable().c_str(), assignReq.managerEpoch, @@ -2345,18 +2390,18 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, assignReq.continueAssignment ? "T" : "F"); } - if (self.managerEpochOk(assignReq.managerEpoch)) { - addActor.send(handleRangeAssign(&self, assignReq, false)); + if (self->managerEpochOk(assignReq.managerEpoch)) { + self->addActor.send(handleRangeAssign(self, assignReq, false)); } else { assignReq.reply.send(AssignBlobRangeReply(false)); } } when(RevokeBlobRangeRequest _req = waitNext(bwInterf.revokeBlobRangeRequest.getFuture())) { state RevokeBlobRangeRequest revokeReq = _req; - --self.stats.numRangesAssigned; + --self->stats.numRangesAssigned; if (BW_DEBUG) { printf("Worker %s revoked range [%s - %s) @ (%lld, %lld):\n dispose=%s\n", - self.id.toString().c_str(), + self->id.toString().c_str(), revokeReq.keyRange.begin.printable().c_str(), revokeReq.keyRange.end.printable().c_str(), revokeReq.managerEpoch, @@ -2364,30 +2409,37 @@ ACTOR Future blobWorker(BlobWorkerInterface bwInterf, revokeReq.dispose ? "T" : "F"); } - if (self.managerEpochOk(revokeReq.managerEpoch)) { - addActor.send(handleRangeRevoke(&self, revokeReq)); + if (self->managerEpochOk(revokeReq.managerEpoch)) { + self->addActor.send(handleRangeRevoke(self, revokeReq)); } else { revokeReq.reply.send(AssignBlobRangeReply(false)); } } - when(AssignBlobRangeRequest granuleToReassign = waitNext(self.granuleUpdateErrors.getFuture())) { - addActor.send(handleRangeAssign(&self, granuleToReassign, true)); + when(AssignBlobRangeRequest granuleToReassign = waitNext(self->granuleUpdateErrors.getFuture())) { + self->addActor.send(handleRangeAssign(self, granuleToReassign, true)); + } + when(HaltBlobWorkerRequest req = waitNext(bwInterf.haltBlobWorker.getFuture())) { + req.reply.send(Void()); + if (self->managerEpochOk(req.managerEpoch)) { + TraceEvent("BlobWorkerHalted", bwInterf.id()).detail("ReqID", req.requesterID); + printf("BW %s was halted\n", bwInterf.id().toString().c_str()); + break; + } } when(wait(collection)) { - if (BW_DEBUG) { - printf("BW actor collection returned, exiting\n"); - } + TraceEvent("BlobWorkerActorCollectionError"); ASSERT(false); throw internal_error(); } } } catch (Error& e) { if (BW_DEBUG) { - printf("Blob worker got error %s, exiting\n", e.name()); + printf("Blob worker got error %s. Exiting...\n", e.name()); } - TraceEvent("BlobWorkerDied", self.id).error(e, true); - throw e; + TraceEvent("BlobWorkerDied", self->id).error(e, true); } + + return Void(); } // TODO add unit tests for assign/revoke range, especially version ordering diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 4ca14d230e..c1f911a6d0 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -251,8 +251,10 @@ if(WITH_PYTHON) add_fdb_test(TEST_FILES slow/ApiCorrectness.toml) add_fdb_test(TEST_FILES slow/ApiCorrectnessAtomicRestore.toml) add_fdb_test(TEST_FILES slow/ApiCorrectnessSwitchover.toml) - add_fdb_test(TEST_FILES fast/BlobGranuleCorrectness.toml) - add_fdb_test(TEST_FILES slow/BlobGranuleCorrectnessLarge.toml) + add_fdb_test(TEST_FILES fast/BlobGranuleCorrectness.toml IGNORE) + add_fdb_test(TEST_FILES slow/BlobGranuleCorrectnessLarge.toml IGNORE) + add_fdb_test(TEST_FILES fast/BlobGranuleCorrectnessClean.toml) + add_fdb_test(TEST_FILES slow/BlobGranuleCorrectnessLargeClean.toml) add_fdb_test(TEST_FILES slow/ClogWithRollbacks.toml) add_fdb_test(TEST_FILES slow/CloggedCycleTest.toml) add_fdb_test(TEST_FILES slow/CloggedStorefront.toml) diff --git a/tests/fast/BlobGranuleCorrectness.toml b/tests/fast/BlobGranuleCorrectness.toml index 59d6be0364..20446ec66c 100644 --- a/tests/fast/BlobGranuleCorrectness.toml +++ b/tests/fast/BlobGranuleCorrectness.toml @@ -8,3 +8,27 @@ testTitle = 'BlobGranuleCorrectnessTest' [[test.workload]] testName = 'BlobGranuleVerifier' testDuration = 120.0 + + [[test.workload]] + testName = 'RandomClogging' + testDuration = 120.0 + + [[test.workload]] + testName = 'Rollback' + meanDelay = 30.0 + testDuration = 120.0 + + [[test.workload]] + testName = 'Attrition' + machinesToKill = 10 + machinesToLeave = 3 + reboot = true + testDuration = 120.0 + + [[test.workload]] + testName = 'Attrition' + machinesToKill = 10 + machinesToLeave = 3 + reboot = true + testDuration = 120.0 + diff --git a/tests/fast/BlobGranuleCorrectnessClean.toml b/tests/fast/BlobGranuleCorrectnessClean.toml new file mode 100644 index 0000000000..168790dd9d --- /dev/null +++ b/tests/fast/BlobGranuleCorrectnessClean.toml @@ -0,0 +1,10 @@ +[[test]] +testTitle = 'BlobGranuleCorrectnessCleanTest' + + [[test.workload]] + testName = 'WriteDuringRead' + testDuration = 120.0 + + [[test.workload]] + testName = 'BlobGranuleVerifier' + testDuration = 120.0 diff --git a/tests/slow/BlobGranuleCorrectnessLarge.toml b/tests/slow/BlobGranuleCorrectnessLarge.toml index e88f315225..edf400b6d9 100644 --- a/tests/slow/BlobGranuleCorrectnessLarge.toml +++ b/tests/slow/BlobGranuleCorrectnessLarge.toml @@ -1,5 +1,5 @@ [[test]] -testTitle = 'BlobGranuleCorrectnessTestLarge' +testTitle = 'BlobGranuleCorrectnessLargeTest' [[test.workload]] testName = 'ReadWrite' @@ -19,3 +19,27 @@ testTitle = 'BlobGranuleCorrectnessTestLarge' [[test.workload]] testName = 'BlobGranuleVerifier' testDuration = 200.0 + + [[test.workload]] + testName = 'RandomClogging' + testDuration = 200.0 + + [[test.workload]] + testName = 'Rollback' + meanDelay = 30.0 + testDuration = 200.0 + + [[test.workload]] + testName = 'Attrition' + machinesToKill = 10 + machinesToLeave = 3 + reboot = true + testDuration = 200.0 + + [[test.workload]] + testName = 'Attrition' + machinesToKill = 10 + machinesToLeave = 3 + reboot = true + testDuration = 200.0 + diff --git a/tests/slow/BlobGranuleCorrectnessLargeClean.toml b/tests/slow/BlobGranuleCorrectnessLargeClean.toml new file mode 100644 index 0000000000..1a12e6f47f --- /dev/null +++ b/tests/slow/BlobGranuleCorrectnessLargeClean.toml @@ -0,0 +1,21 @@ +[[test]] +testTitle = 'BlobGranuleCorrectnessLargeCleanTest' + + [[test.workload]] + testName = 'ReadWrite' + testDuration = 200.0 + transactionsPerSecond = 200 + 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