diff --git a/fdbcli/fdbcli.actor.cpp b/fdbcli/fdbcli.actor.cpp index 608324cd98..3265126bd4 100644 --- a/fdbcli/fdbcli.actor.cpp +++ b/fdbcli/fdbcli.actor.cpp @@ -3580,12 +3580,32 @@ ACTOR Future cli(CLIOptions opt, LineNoise* plinenoise) { } } } else if (tokencmp(tokens[1], "get")) { - if (tokens.size() != 3) { + if (tokens.size() < 3 || tokens.size() > 5) { printUsage(tokens[0]); is_error = true; continue; } - Standalone> res = wait(db->getRangeFeedMutations(tokens[2])); + Version begin = 0; + Version end = std::numeric_limits::max(); + if (tokens.size() > 3) { + int n = 0; + if (sscanf(tokens[3].toString().c_str(), "%ld%n", &begin, &n) != 1 || + n != tokens[3].size()) { + printUsage(tokens[0]); + is_error = true; + continue; + } + } + if (tokens.size() > 4) { + int n = 0; + if (sscanf(tokens[4].toString().c_str(), "%ld%n", &end, &n) != 1 || n != tokens[4].size()) { + printUsage(tokens[0]); + is_error = true; + continue; + } + } + Standalone> res = + wait(db->getRangeFeedMutations(tokens[2], begin, end)); for (auto& it : res) { for (auto& it2 : it.mutations) { printf("%lld %s\n", it.version, it2.toString().c_str()); diff --git a/fdbclient/DatabaseContext.h b/fdbclient/DatabaseContext.h index e6d98cef6e..8a7fe4021f 100644 --- a/fdbclient/DatabaseContext.h +++ b/fdbclient/DatabaseContext.h @@ -252,8 +252,11 @@ public: // Management API, create snapshot Future createSnapshot(StringRef uid, StringRef snapshot_command); - Future>> getRangeFeedMutations(StringRef rangeID, - KeyRangeRef range = allKeys); + Future>> getRangeFeedMutations( + StringRef rangeID, + Version begin = 0, + Version end = std::numeric_limits::max(), + KeyRange range = allKeys); Future>> getOverlappingRangeFeeds(KeyRangeRef ranges, Version minVersion); Future popRangeFeedMutations(StringRef rangeID, Version version); diff --git a/fdbclient/NativeAPI.actor.cpp b/fdbclient/NativeAPI.actor.cpp index c198994f48..09b8070cc2 100644 --- a/fdbclient/NativeAPI.actor.cpp +++ b/fdbclient/NativeAPI.actor.cpp @@ -6519,7 +6519,9 @@ Future DatabaseContext::createSnapshot(StringRef uid, StringRef snapshot_c ACTOR Future>> getRangeFeedMutationsActor(Reference db, StringRef rangeID, - KeyRangeRef range) { + Version begin, + Version end, + KeyRange range) { state Database cx(db); state Transaction tr(cx); state Key rangeIDKey = rangeID.withPrefix(rangeFeedPrefix); @@ -6543,6 +6545,8 @@ ACTOR Future>> getRangeFeedMutation state RangeFeedRequest req; req.rangeID = rangeID; + req.begin = begin; + req.end = end; RangeFeedReply rep = wait(loadBalance(cx.getPtr(), locations[0].second, @@ -6555,8 +6559,10 @@ ACTOR Future>> getRangeFeedMutation } Future>> DatabaseContext::getRangeFeedMutations(StringRef rangeID, - KeyRangeRef range) { - return getRangeFeedMutationsActor(Reference::addRef(this), rangeID, range); + Version begin, + Version end, + KeyRange range) { + return getRangeFeedMutationsActor(Reference::addRef(this), rangeID, begin, end, range); } ACTOR Future>> getOverlappingRangeFeedsActor(Reference db, diff --git a/fdbclient/StorageServerInterface.h b/fdbclient/StorageServerInterface.h index 9efea0e732..7661398a58 100644 --- a/fdbclient/StorageServerInterface.h +++ b/fdbclient/StorageServerInterface.h @@ -663,6 +663,8 @@ struct RangeFeedReply { struct RangeFeedRequest { constexpr static FileIdentifier file_identifier = 10726174; Key rangeID; + Version begin = 0; + Version end = 0; ReplyPromise reply; RangeFeedRequest() {} @@ -670,7 +672,7 @@ struct RangeFeedRequest { template void serialize(Ar& ar) { - serializer(ar, rangeID, reply); + serializer(ar, rangeID, begin, end, reply); } }; diff --git a/fdbserver/storageserver.actor.cpp b/fdbserver/storageserver.actor.cpp index 3bf85cfd93..b6827d00d8 100644 --- a/fdbserver/storageserver.actor.cpp +++ b/fdbserver/storageserver.actor.cpp @@ -313,6 +313,7 @@ struct FetchInjectionInfo { struct RangeFeedInfo : ReferenceCounted { std::deque> mutations; Version durableVersion = invalidVersion; + Version emptyVersion = 0; KeyRange range; Key id; }; @@ -1531,7 +1532,8 @@ ACTOR Future getRangeFeedMutations(StorageServer* data, RangeFee state RangeFeedReply reply; wait(delay(0)); auto& feedInfo = data->uidRangeFeed[req.rangeID]; - if (feedInfo->durableVersion == invalidVersion) { + if (req.end <= feedInfo->emptyVersion + 1) { + } else if (feedInfo->durableVersion == invalidVersion || req.begin > feedInfo->durableVersion) { for (auto& it : data->uidRangeFeed[req.rangeID]->mutations) { reply.mutations.push_back(reply.arena, it); } @@ -1539,10 +1541,7 @@ ACTOR Future getRangeFeedMutations(StorageServer* data, RangeFee state std::deque> mutationsDeque = data->uidRangeFeed[req.rangeID]->mutations; RangeResult res = wait(data->storage.readRange( - KeyRangeRef(rangeFeedDurableKey(req.rangeID, 0), rangeFeedDurableKey(req.rangeID, data->version.get())))); - if (res.empty()) { - data->uidRangeFeed[req.rangeID]->durableVersion = invalidVersion; - } + KeyRangeRef(rangeFeedDurableKey(req.rangeID, req.begin), rangeFeedDurableKey(req.rangeID, req.end)))); Version lastVersion = invalidVersion; for (auto& kv : res) { Key id; @@ -1553,10 +1552,35 @@ ACTOR Future getRangeFeedMutations(StorageServer* data, RangeFee lastVersion = version; } for (auto& it : mutationsDeque) { + if (it.version >= req.end) { + break; + } if (it.version > lastVersion) { reply.mutations.push_back(reply.arena, it); } } + if (res.empty()) { + auto& feedInfo = data->uidRangeFeed[req.rangeID]; + if (req.end > feedInfo->durableVersion) { + if (req.begin == 0) { + feedInfo->durableVersion = req.end > data->storageVersion() ? invalidVersion : req.end; + } else { + RangeResult emp = wait(data->storage.readRange( + KeyRangeRef(rangeFeedDurableKey(req.rangeID, 0), rangeFeedDurableKey(req.rangeID, req.end)), + -1)); + + auto& feedInfo = data->uidRangeFeed[req.rangeID]; + if (emp.empty()) { + feedInfo->durableVersion = req.end > data->storageVersion() ? invalidVersion : req.end; + } else { + Key id; + Version version; + std::tie(id, version) = decodeRangeFeedDurableKey(emp[0].key); + feedInfo->durableVersion = version; + } + } + } + } } return reply; } @@ -2945,7 +2969,7 @@ static const KeyRangeRef persistRangeFeedKeys = KeyRangeRef(LiteralStringRef(PERSIST_PREFIX "RF/"), LiteralStringRef(PERSIST_PREFIX "RF0")); // data keys are unmangled (but never start with PERSIST_PREFIX because they are always in allKeys) -ACTOR Future fetchRangeFeed(StorageServer* data, Key rangeId, KeyRange range) { +ACTOR Future fetchRangeFeed(StorageServer* data, Key rangeId, KeyRange range, Version fetchVersion) { TraceEvent("FetchRangeFeed", data->thisServerID) .detail("RangeID", rangeId.printable()) @@ -2966,9 +2990,11 @@ ACTOR Future fetchRangeFeed(StorageServer* data, Key rangeId, KeyRange ran rangeFeedValue(range))); state Standalone> mutations = - wait(data->cx->getRangeFeedMutations(rangeId, range)); + wait(data->cx->getRangeFeedMutations(rangeId, 0, fetchVersion, range)); state RangeFeedRequest req; req.rangeID = rangeId; + req.begin = 0; + req.end = fetchVersion; state RangeFeedReply rep = wait(getRangeFeedMutations(data, req)); state int mLoc = 0; state int rLoc = 0; @@ -3000,7 +3026,7 @@ ACTOR Future dispatchRangeFeeds(StorageServer* data, UID fetchKeysID, KeyR state std::vector> feeds = wait(data->cx->getOverlappingRangeFeeds(keys, fetchVersion)); for (auto& feed : feeds) { - feedFetches[feed.first] = fetchRangeFeed(data, feed.first, feed.second); + feedFetches[feed.first] = fetchRangeFeed(data, feed.first, feed.second, fetchVersion); } loop { @@ -3752,6 +3778,7 @@ private: Reference rangeFeedInfo(new RangeFeedInfo()); rangeFeedInfo->range = rangeFeedRange; rangeFeedInfo->id = rangeFeedId; + rangeFeedInfo->emptyVersion = currentVersion - 1; data->uidRangeFeed[rangeFeedId] = rangeFeedInfo; auto rs = data->keyRangeFeed.modify(rangeFeedRange); for (auto r = rs.begin(); r != rs.end(); ++r) { @@ -4277,11 +4304,13 @@ ACTOR Future updateStorage(StorageServer* data) { state int curFeed = 0; while (curFeed < updatedRangeFeeds.size()) { auto info = data->uidRangeFeed[updatedRangeFeeds[curFeed]]; - while (info->mutations.front().version < newOldestVersion) { + for (auto& it : info->mutations) { + if (it.version >= newOldestVersion) { + break; + } data->storage.writeKeyValue(KeyValueRef(rangeFeedDurableKey(info->id, info->mutations.front().version), rangeFeedDurableValue(info->mutations.front().mutations))); info->durableVersion = info->mutations.front().version; - info->mutations.pop_front(); } wait(yield(TaskPriority::UpdateStorage)); curFeed++; @@ -4320,6 +4349,16 @@ ACTOR Future updateStorage(StorageServer* data) { throw please_reboot(); } + curFeed = 0; + while (curFeed < updatedRangeFeeds.size()) { + auto info = data->uidRangeFeed[updatedRangeFeeds[curFeed]]; + while (info->mutations.front().version < newOldestVersion) { + info->mutations.pop_front(); + } + wait(yield(TaskPriority::UpdateStorage)); + curFeed++; + } + durableInProgress.send(Void()); wait(delay(0, TaskPriority::UpdateStorage)); // Setting durableInProgess could cause the storage server to shut // down, so delay to check for cancellation @@ -5240,14 +5279,17 @@ ACTOR Future serveRangeFeedPopRequests(StorageServer* self, FutureStreamuidRangeFeed[req.rangeID]; - while (!feed->mutations.empty() && feed->mutations.front().version < req.version) { - self->uidRangeFeed[req.rangeID]->mutations.pop_front(); - } - if (feed->durableVersion != invalidVersion) { - self->storage.clearRange( - KeyRangeRef(rangeFeedDurableKey(feed->id, 0), rangeFeedDurableKey(feed->id, req.version))); - if (req.version > feed->durableVersion) { - feed->durableVersion = invalidVersion; + if (req.version - 1 > feed->emptyVersion) { + feed->emptyVersion = req.version - 1; + while (!feed->mutations.empty() && feed->mutations.front().version < req.version) { + self->uidRangeFeed[req.rangeID]->mutations.pop_front(); + } + if (feed->durableVersion != invalidVersion) { + self->storage.clearRange( + KeyRangeRef(rangeFeedDurableKey(feed->id, 0), rangeFeedDurableKey(feed->id, req.version))); + if (req.version > feed->durableVersion) { + feed->durableVersion = invalidVersion; + } } } TraceEvent("RangeFeedPopQuery", self->thisServerID)