1923 lines
68 KiB
C++
1923 lines
68 KiB
C++
/*
|
|
* LogSystemPeekCursor.cpp
|
|
*
|
|
* This source file is part of the FoundationDB open source project
|
|
*
|
|
* Copyright 2013-2026 Apple Inc. and the FoundationDB project authors
|
|
*
|
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
|
* you may not use this file except in compliance with the License.
|
|
* You may obtain a copy of the License at
|
|
*
|
|
* http://www.apache.org/licenses/LICENSE-2.0
|
|
*
|
|
* Unless required by applicable law or agreed to in writing, software
|
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
* See the License for the specific language governing permissions and
|
|
* limitations under the License.
|
|
*/
|
|
|
|
#include "fdbserver/logsystem/LogSystem.h"
|
|
#include "LogSystemTypes.h"
|
|
#include "fdbrpc/FailureMonitor.h"
|
|
#include "fdbserver/core/Knobs.h"
|
|
#include "fdbserver/core/MutationTracking.h"
|
|
#include "fdbrpc/ReplicationUtils.h"
|
|
#include "flow/DebugTrace.h"
|
|
#include "flow/CoroUtils.h"
|
|
#include "flow/UnitTest.h"
|
|
|
|
// Returns a timeout future for peek reply detection. Used by both the non-parallel
|
|
// (serverPeekGetMoreImpl) and parallel (serverPeekParallelGetMoreImpl) peek paths
|
|
// to detect dead/stale TLog endpoints that silently drop requests. When the timeout
|
|
// fires, the caller re-sends the peek. Set PEEK_REPLY_TIMEOUT to 0 to disable.
|
|
static Future<Void> peekReplyTimeout(ServerPeekCursor const* self) {
|
|
return (self->interf->get().present() && SERVER_KNOBS->PEEK_REPLY_TIMEOUT > 0)
|
|
? delay(SERVER_KNOBS->PEEK_REPLY_TIMEOUT)
|
|
: Never();
|
|
}
|
|
|
|
Future<Void> tryEstablishPeekStreamImpl(ServerPeekCursor* self) {
|
|
co_await IFailureMonitor::failureMonitor().onStateEqual(
|
|
self->interf->get().interf().peekStreamMessages.getEndpoint(), FailureStatus(false));
|
|
|
|
auto req = TLogPeekStreamRequest(self->messageVersion.version,
|
|
self->tag,
|
|
self->returnIfBlocked,
|
|
std::numeric_limits<int>::max(),
|
|
self->end.version,
|
|
self->returnEmptyIfStopped);
|
|
self->peekReplyStream = self->interf->get().interf().peekStreamMessages.getReplyStream(req);
|
|
DebugLogTraceEvent(SevDebug, "SPC_StreamCreated", self->randomID)
|
|
.detail("Tag", self->tag)
|
|
.detail("PeerAddr", self->interf->get().interf().peekStreamMessages.getEndpoint().getPrimaryAddress())
|
|
.detail("PeerAddress", self->interf->get().interf().peekStreamMessages.getEndpoint().getPrimaryAddress())
|
|
.detail("PeerToken", self->interf->get().interf().peekStreamMessages.getEndpoint().token);
|
|
}
|
|
|
|
// create a peek stream for cursor when it's possible
|
|
Future<Void> tryEstablishPeekStream(ServerPeekCursor* self) {
|
|
if (self->peekReplyStream.present())
|
|
return Void();
|
|
else if (!self->interf || !self->interf->get().present()) {
|
|
self->peekReplyStream.reset();
|
|
return Never();
|
|
}
|
|
return tryEstablishPeekStreamImpl(self);
|
|
}
|
|
|
|
ServerPeekCursor::ServerPeekCursor(Reference<AsyncVar<OptionalInterface<TLogInterface>>> const& interf,
|
|
Tag tag,
|
|
Version begin,
|
|
Version end,
|
|
bool returnIfBlocked,
|
|
bool parallelGetMore,
|
|
bool returnEmptyIfStopped)
|
|
: interf(interf), tag(tag), rd(results.arena, results.messages, Unversioned()), messageVersion(begin), end(end),
|
|
poppedVersion(0), hasMsg(false), randomID(deterministicRandom()->randomUniqueID()),
|
|
returnIfBlocked(returnIfBlocked), onlySpilled(false), parallelGetMore(parallelGetMore),
|
|
usePeekStream(SERVER_KNOBS->PEEK_USING_STREAMING), sequence(0), lastReset(0), resetCheck(Void()), slowReplies(0),
|
|
fastReplies(0), unknownReplies(0), returnEmptyIfStopped(returnEmptyIfStopped), replyByteLimit(0) {
|
|
this->results.maxKnownVersion = 0;
|
|
this->results.minKnownCommittedVersion = 0;
|
|
DebugLogTraceEvent(SevDebug, "SPC_Starting", randomID)
|
|
.detail("Tag", tag.toString())
|
|
.detail("Parallel", parallelGetMore)
|
|
.detail("Interf", interf && interf->get().present() ? interf->get().id() : UID())
|
|
.detail("UsePeekStream", usePeekStream)
|
|
.detail("Begin", begin)
|
|
.detail("End", end);
|
|
}
|
|
|
|
ServerPeekCursor::ServerPeekCursor(TLogPeekReply const& results,
|
|
LogMessageVersion const& messageVersion,
|
|
LogMessageVersion const& end,
|
|
TagsAndMessage const& message,
|
|
bool hasMsg,
|
|
Version poppedVersion,
|
|
Tag tag)
|
|
: tag(tag), results(results), rd(results.arena, results.messages, Unversioned()), messageVersion(messageVersion),
|
|
end(end), poppedVersion(poppedVersion), messageAndTags(message), hasMsg(hasMsg),
|
|
randomID(deterministicRandom()->randomUniqueID()), returnIfBlocked(false), onlySpilled(false),
|
|
parallelGetMore(false), usePeekStream(false), sequence(0), lastReset(0), resetCheck(Void()), slowReplies(0),
|
|
fastReplies(0), unknownReplies(0), returnEmptyIfStopped(false), replyByteLimit(0) {
|
|
//TraceEvent("SPC_Clone", randomID);
|
|
this->results.maxKnownVersion = 0;
|
|
this->results.minKnownCommittedVersion = 0;
|
|
if (hasMsg)
|
|
nextMessage();
|
|
|
|
advanceTo(messageVersion);
|
|
}
|
|
|
|
Reference<ServerPeekCursor> ServerPeekCursor::cloneServerNoMore() {
|
|
return makeReference<ServerPeekCursor>(results, messageVersion, end, messageAndTags, hasMsg, poppedVersion, tag);
|
|
}
|
|
|
|
Reference<IReplayPeekCursor> ServerPeekCursor::cloneNoMore() {
|
|
return cloneServerNoMore();
|
|
}
|
|
|
|
void ServerPeekCursor::setProtocolVersion(ProtocolVersion version) {
|
|
rd.setProtocolVersion(version);
|
|
}
|
|
|
|
Arena& ServerPeekCursor::arena() {
|
|
return results.arena;
|
|
}
|
|
|
|
ArenaReader* ServerPeekCursor::reader() {
|
|
return &rd;
|
|
}
|
|
|
|
bool ServerPeekCursor::hasMessage() const {
|
|
//TraceEvent("SPC_HasMessage", randomID).detail("HasMsg", hasMsg);
|
|
return hasMsg;
|
|
}
|
|
|
|
void ServerPeekCursor::nextMessage() {
|
|
DebugLogTraceEvent("SPC_NextMessage", randomID)
|
|
.detail("Tag", tag.toString())
|
|
.detail("MessageVersion", messageVersion.toString());
|
|
ASSERT(hasMsg);
|
|
if (rd.empty()) {
|
|
messageVersion.reset(std::min(results.end, end.version));
|
|
hasMsg = false;
|
|
return;
|
|
}
|
|
if (*(int32_t*)rd.peekBytes(4) == VERSION_HEADER) {
|
|
// A version
|
|
int32_t dummy;
|
|
Version ver;
|
|
rd >> dummy >> ver;
|
|
|
|
//TraceEvent("SPC_ProcessSeq", randomID).detail("MessageVersion", messageVersion.toString()).detail("Ver", ver).detail("Tag", tag.toString());
|
|
// ASSERT( ver >= messageVersion.version );
|
|
|
|
messageVersion.reset(ver);
|
|
|
|
if (messageVersion >= end) {
|
|
messageVersion = end;
|
|
hasMsg = false;
|
|
return;
|
|
}
|
|
ASSERT(!rd.empty());
|
|
}
|
|
|
|
messageAndTags.loadFromArena(&rd, &messageVersion.sub);
|
|
DEBUG_TAGS_AND_MESSAGE("ServerPeekCursor", messageVersion.version, messageAndTags.getRawMessage(), this->randomID);
|
|
// Rewind and consume the header so that reader() starts from the message.
|
|
rd.rewind();
|
|
rd.readBytes(TagsAndMessage::getHeaderSize(messageAndTags.tags.size()));
|
|
hasMsg = true;
|
|
DebugLogTraceEvent("SPC_NextMessageB", randomID)
|
|
.detail("Tag", tag.toString())
|
|
.detail("MessageVersion", messageVersion.toString());
|
|
}
|
|
|
|
StringRef ServerPeekCursor::getMessage() {
|
|
DebugLogTraceEvent("SPC_GetMessage", randomID).detail("Tag", tag.toString());
|
|
StringRef message = messageAndTags.getMessageWithoutTags();
|
|
rd.readBytes(message.size()); // Consumes the message.
|
|
return message;
|
|
}
|
|
|
|
StringRef ServerPeekCursor::getMessageWithTags() {
|
|
StringRef rawMessage = messageAndTags.getRawMessage();
|
|
rd.readBytes(rawMessage.size() -
|
|
TagsAndMessage::getHeaderSize(messageAndTags.tags.size())); // Consumes the message.
|
|
return rawMessage;
|
|
}
|
|
|
|
VectorRef<Tag> ServerPeekCursor::getTags() const {
|
|
return messageAndTags.tags;
|
|
}
|
|
|
|
void ServerPeekCursor::advanceTo(LogMessageVersion n) {
|
|
//TraceEvent("SPC_AdvanceTo", randomID).detail("N", n.toString());
|
|
while (messageVersion < n && hasMessage()) {
|
|
getMessage();
|
|
nextMessage();
|
|
}
|
|
|
|
if (hasMessage())
|
|
return;
|
|
|
|
// if( more.isValid() && !more.isReady() ) more.cancel();
|
|
|
|
if (messageVersion < n) {
|
|
messageVersion = n;
|
|
}
|
|
}
|
|
|
|
// This function is called after the cursor received one TLogPeekReply to update its members, which is the common logic
|
|
// in getMore helper functions.
|
|
void updateCursorWithReply(ServerPeekCursor* self, const TLogPeekReply& res) {
|
|
self->results = res;
|
|
self->onlySpilled = res.onlySpilled;
|
|
if (res.popped.present())
|
|
self->poppedVersion = std::min(std::max(self->poppedVersion, res.popped.get()), self->end.version);
|
|
self->rd = ArenaReader(self->results.arena, self->results.messages, Unversioned());
|
|
LogMessageVersion skipSeq = self->messageVersion;
|
|
self->hasMsg = true;
|
|
self->nextMessage();
|
|
self->advanceTo(skipSeq);
|
|
}
|
|
|
|
Future<Void> resetChecker(ServerPeekCursor* self, NetworkAddress addr) {
|
|
self->slowReplies = 0;
|
|
self->unknownReplies = 0;
|
|
self->fastReplies = 0;
|
|
co_await delay(SERVER_KNOBS->PEEK_STATS_INTERVAL);
|
|
TraceEvent("SlowPeekStats", self->randomID)
|
|
.detail("PeerAddress", addr)
|
|
.detail("SlowReplies", self->slowReplies)
|
|
.detail("FastReplies", self->fastReplies)
|
|
.detail("UnknownReplies", self->unknownReplies);
|
|
|
|
if (self->slowReplies >= SERVER_KNOBS->PEEK_STATS_SLOW_AMOUNT &&
|
|
self->slowReplies / double(self->slowReplies + self->fastReplies) >= SERVER_KNOBS->PEEK_STATS_SLOW_RATIO) {
|
|
|
|
TraceEvent("ConnectionResetSlowPeek", self->randomID)
|
|
.detail("PeerAddress", addr)
|
|
.detail("SlowReplies", self->slowReplies)
|
|
.detail("FastReplies", self->fastReplies)
|
|
.detail("UnknownReplies", self->unknownReplies);
|
|
FlowTransport::transport().resetConnection(addr);
|
|
self->lastReset = now();
|
|
}
|
|
}
|
|
|
|
Future<TLogPeekReply> recordRequestMetrics(ServerPeekCursor* self, NetworkAddress addr, Future<TLogPeekReply> in) {
|
|
Error err;
|
|
try {
|
|
double startTime = now();
|
|
TLogPeekReply t = co_await in;
|
|
if (now() - self->lastReset > SERVER_KNOBS->PEEK_RESET_INTERVAL) {
|
|
if (now() - startTime > SERVER_KNOBS->PEEK_MAX_LATENCY) {
|
|
if (t.messages.size() >= SERVER_KNOBS->DESIRED_TOTAL_BYTES || SERVER_KNOBS->PEEK_COUNT_SMALL_MESSAGES) {
|
|
if (self->resetCheck.isReady()) {
|
|
self->resetCheck = resetChecker(self, addr);
|
|
}
|
|
self->slowReplies++;
|
|
} else {
|
|
self->unknownReplies++;
|
|
}
|
|
} else {
|
|
self->fastReplies++;
|
|
}
|
|
}
|
|
co_return t;
|
|
} catch (Error& e) {
|
|
err = e;
|
|
}
|
|
if (err.code() != error_code_broken_promise)
|
|
throw err;
|
|
co_await Future<Void>(Never()); // never return
|
|
throw internal_error(); // does not happen
|
|
}
|
|
|
|
Future<Void> serverPeekParallelGetMoreImpl(ServerPeekCursor* self, TaskPriority taskID) {
|
|
while (true) {
|
|
DebugLogTraceEvent("SPC_GetMoreP", self->randomID)
|
|
.detail("Tag", self->tag.toString())
|
|
.detail("Has", self->hasMessage())
|
|
.detail("Begin", self->messageVersion.version)
|
|
.detail("Parallel", self->parallelGetMore)
|
|
.detail("Seq", self->sequence)
|
|
.detail("Sizes", self->futureResults.size())
|
|
.detail("Interf", self->interf->get().present() ? self->interf->get().id() : UID());
|
|
|
|
Version expectedBegin = self->messageVersion.version;
|
|
try {
|
|
if (self->parallelGetMore || self->onlySpilled) {
|
|
while (self->futureResults.size() < SERVER_KNOBS->PARALLEL_GET_MORE_REQUESTS &&
|
|
self->interf->get().present()) {
|
|
self->futureResults.push_back(recordRequestMetrics(
|
|
self,
|
|
self->interf->get().interf().peekMessages.getEndpoint().getPrimaryAddress(),
|
|
self->interf->get().interf().peekMessages.getReply(
|
|
TLogPeekRequest(self->messageVersion.version,
|
|
self->tag,
|
|
self->returnIfBlocked,
|
|
self->onlySpilled,
|
|
std::make_pair(self->randomID, self->sequence++),
|
|
self->end.version,
|
|
self->returnEmptyIfStopped,
|
|
self->replyByteLimit),
|
|
taskID)));
|
|
}
|
|
if (self->sequence == std::numeric_limits<decltype(self->sequence)>::max()) {
|
|
throw operation_obsolete();
|
|
}
|
|
} else if (self->futureResults.empty()) {
|
|
co_return;
|
|
}
|
|
|
|
if (self->hasMessage())
|
|
co_return;
|
|
|
|
Future<TLogPeekReply> peekReply = self->interf->get().present() ? self->futureResults.front() : Never();
|
|
// See peekReplyTimeout() and serverPeekGetMoreImpl for the non-parallel equivalent.
|
|
auto res = co_await race(peekReply, self->interfaceChanged, peekReplyTimeout(self));
|
|
if (res.index() == 0) {
|
|
TLogPeekReply reply = std::get<0>(std::move(res));
|
|
|
|
if (reply.begin.get() != expectedBegin) {
|
|
throw operation_obsolete();
|
|
}
|
|
self->futureResults.pop_front();
|
|
updateCursorWithReply(self, reply);
|
|
DebugLogTraceEvent("SPC_GetMoreReply", self->randomID)
|
|
.detail("Has", self->hasMessage())
|
|
.detail("Tag", self->tag.toString())
|
|
.detail("End", reply.end)
|
|
.detail("Size", self->futureResults.size())
|
|
.detail("Popped", reply.popped.present() ? reply.popped.get() : 0);
|
|
co_return;
|
|
} else if (res.index() == 1) {
|
|
self->interfaceChanged = self->interf->onChange();
|
|
self->randomID = deterministicRandom()->randomUniqueID();
|
|
self->sequence = 0;
|
|
self->onlySpilled = false;
|
|
self->futureResults.clear();
|
|
} else if (res.index() == 2) {
|
|
// Timeout fired — no reply within PEEK_REPLY_TIMEOUT. Re-send the peek.
|
|
// This handles dead/stale TLog endpoints that silently drop requests.
|
|
DebugLogTraceEvent("PeekParallelReplyTimeout", self->randomID)
|
|
.detail("Tag", self->tag.toString())
|
|
.detail("Version", self->messageVersion.version);
|
|
self->randomID = deterministicRandom()->randomUniqueID();
|
|
self->sequence = 0;
|
|
self->futureResults.clear();
|
|
} else {
|
|
UNREACHABLE();
|
|
}
|
|
} catch (Error& e) {
|
|
DebugLogTraceEvent("PeekCursorError", self->randomID)
|
|
.error(e)
|
|
.detail("Tag", self->tag.toString())
|
|
.detail("Begin", self->messageVersion.version)
|
|
.detail("Interf", self->interf->get().present() ? self->interf->get().id() : UID());
|
|
|
|
if (e.code() == error_code_end_of_stream) {
|
|
self->end.reset(self->messageVersion.version);
|
|
co_return;
|
|
} else if (e.code() == error_code_timed_out || e.code() == error_code_operation_obsolete) {
|
|
TraceEvent ev("PeekCursorTimedOut", self->randomID);
|
|
// We *should* never get timed_out(), as it means the TLog got stuck while handling a parallel peek,
|
|
// and thus we've likely just wasted 10min.
|
|
// timed_out() is sent by cleanupPeekTrackers as value PEEK_TRACKER_EXPIRATION_TIME
|
|
//
|
|
// A cursor for a log router can be delayed indefinitely during a network partition, so only fail
|
|
// simulation tests sufficiently far after we finish simulating network partitions.
|
|
CODE_PROBE(e.code() == error_code_timed_out, "peek cursor timed out", probe::decoration::rare);
|
|
if (g_network->isSimulated() && now() >= g_simulator->connectionFailureEnableTime +
|
|
FLOW_KNOBS->SIM_SPEEDUP_AFTER_SECONDS +
|
|
SERVER_KNOBS->PEEK_TRACKER_EXPIRATION_TIME) {
|
|
ASSERT_WE_THINK(e.code() == error_code_operation_obsolete ||
|
|
SERVER_KNOBS->PEEK_TRACKER_EXPIRATION_TIME < 10);
|
|
}
|
|
self->interfaceChanged = self->interf->onChange();
|
|
self->randomID = deterministicRandom()->randomUniqueID();
|
|
self->sequence = 0;
|
|
self->futureResults.clear();
|
|
ev.error(e)
|
|
.detail("Tag", self->tag.toString())
|
|
.detail("Begin", self->messageVersion.version)
|
|
.detail("NewID", self->randomID)
|
|
.detail("Interf", self->interf->get().present() ? self->interf->get().id() : UID());
|
|
} else {
|
|
throw e;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
Future<Void> serverPeekParallelGetMore(ServerPeekCursor* self, TaskPriority taskID) {
|
|
if (!self->interf || self->isExhausted()) {
|
|
if (self->hasMessage())
|
|
return Void();
|
|
return Never();
|
|
}
|
|
|
|
if (!self->interfaceChanged.isValid()) {
|
|
self->interfaceChanged = self->interf->onChange();
|
|
}
|
|
|
|
return serverPeekParallelGetMoreImpl(self, taskID);
|
|
}
|
|
|
|
Future<Void> serverPeekStreamGetMoreImpl(ServerPeekCursor* self, TaskPriority taskID) {
|
|
while (true) {
|
|
Optional<Error> err;
|
|
try {
|
|
Version expectedBegin = self->messageVersion.version;
|
|
Future<TLogPeekReply> fPeekReply = self->peekReplyStream.present()
|
|
? map(waitAndForward(self->peekReplyStream.get().getFuture()),
|
|
[](const TLogPeekStreamReply& r) { return r.rep; })
|
|
: Never();
|
|
Future<Void> establishStream = self->peekReplyStream.present() ? Never() : tryEstablishPeekStream(self);
|
|
Future<TLogPeekReply> metricsReply =
|
|
self->peekReplyStream.present()
|
|
? recordRequestMetrics(
|
|
self,
|
|
self->interf->get().interf().peekStreamMessages.getEndpoint().getPrimaryAddress(),
|
|
fPeekReply)
|
|
: Never();
|
|
auto res = co_await race(establishStream, self->interf->onChange(), metricsReply);
|
|
if (res.index() == 0) {
|
|
} else if (res.index() == 1) {
|
|
self->onlySpilled = false;
|
|
self->peekReplyStream.reset();
|
|
} else if (res.index() == 2) {
|
|
TLogPeekReply reply = std::get<2>(std::move(res));
|
|
|
|
if (reply.begin.get() != expectedBegin) {
|
|
throw operation_obsolete();
|
|
}
|
|
updateCursorWithReply(self, reply);
|
|
DebugLogTraceEvent(SevDebug, "SPC_GetMoreB", self->randomID)
|
|
.detail("Tag", self->tag)
|
|
.detail("Has", self->hasMessage())
|
|
.detail("End", reply.end)
|
|
.detail("Popped", reply.popped.present() ? reply.popped.get() : 0);
|
|
|
|
// NOTE: delay is necessary here since ReplyPromiseStream delivers reply on high priority. Here we
|
|
// change the priority to the intended one.
|
|
co_await delay(0, taskID);
|
|
co_return;
|
|
} else {
|
|
UNREACHABLE();
|
|
}
|
|
} catch (Error& e) {
|
|
err = e;
|
|
}
|
|
if (err.present()) {
|
|
DebugLogTraceEvent(SevDebug, "SPC_GetMoreB_Error", self->randomID)
|
|
.errorUnsuppressed(err.get())
|
|
.detail("Tag", self->tag);
|
|
if (err.get().code() == error_code_connection_failed || err.get().code() == error_code_operation_obsolete ||
|
|
err.get().code() == error_code_request_maybe_delivered) {
|
|
// NOTE: delay in order to avoid the endless retry loop block other tasks
|
|
self->peekReplyStream.reset();
|
|
co_await delay(0);
|
|
} else if (err.get().code() == error_code_end_of_stream) {
|
|
self->peekReplyStream.reset();
|
|
self->end.reset(self->messageVersion.version);
|
|
co_return;
|
|
} else {
|
|
throw err.get();
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
Future<Void> serverPeekStreamGetMore(ServerPeekCursor* self, TaskPriority taskID) {
|
|
if (!self->interf || self->isExhausted()) {
|
|
self->peekReplyStream.reset();
|
|
if (self->hasMessage())
|
|
return Void();
|
|
return Never();
|
|
}
|
|
return serverPeekStreamGetMoreImpl(self, taskID);
|
|
}
|
|
|
|
Future<Void> serverPeekGetMoreImpl(ServerPeekCursor* self, TaskPriority taskID) {
|
|
try {
|
|
while (true) {
|
|
Future<TLogPeekReply> peekReply = Never();
|
|
if (self->interf->get().present()) {
|
|
peekReply = brokenPromiseToNever(
|
|
self->interf->get().interf().peekMessages.getReply(TLogPeekRequest(self->messageVersion.version,
|
|
self->tag,
|
|
self->returnIfBlocked,
|
|
self->onlySpilled,
|
|
Optional<std::pair<UID, int>>(),
|
|
self->end.version,
|
|
self->returnEmptyIfStopped,
|
|
self->replyByteLimit),
|
|
taskID));
|
|
}
|
|
// Race between: (1) peek reply, (2) interface change, and optionally
|
|
// (3) sender-side timeout — see peekReplyTimeout() and
|
|
// serverPeekParallelGetMoreImpl for the parallel equivalent.
|
|
auto res = co_await race(peekReply, self->interf->onChange(), peekReplyTimeout(self));
|
|
if (res.index() == 0) {
|
|
TLogPeekReply reply = std::get<0>(std::move(res));
|
|
|
|
updateCursorWithReply(self, reply);
|
|
DebugLogTraceEvent("SPC_GetMoreB", self->randomID)
|
|
.detail("Tag", self->tag.toString())
|
|
.detail("Has", self->hasMessage())
|
|
.detail("End", reply.end)
|
|
.detail("Popped", reply.popped.present() ? reply.popped.get() : 0);
|
|
co_return;
|
|
} else if (res.index() == 1) {
|
|
self->onlySpilled = false;
|
|
} else if (res.index() == 2) {
|
|
// Sender-side timeout fired — no reply received within PEEK_REPLY_TIMEOUT.
|
|
// Don't throw — just retry the peek. This handles both cases:
|
|
// - Dead TLog: the retry will also timeout, but eventually the interface
|
|
// will change (cluster controller detects the failure) and we'll break
|
|
// out via interf->onChange().
|
|
// - Clogged TLog: the retry may succeed if the clog clears, or timeout
|
|
// again harmlessly.
|
|
// The key benefit: this resets the brokenPromiseToNever future, allowing
|
|
// the peek to be re-sent. Without this, a peek to a dead endpoint would
|
|
// wait forever because the original future is permanently stuck.
|
|
DebugLogTraceEvent("PeekReplyTimeout", self->randomID)
|
|
.detail("Tag", self->tag.toString())
|
|
.detail("Version", self->messageVersion.version);
|
|
continue;
|
|
} else {
|
|
UNREACHABLE();
|
|
}
|
|
}
|
|
} catch (Error& e) {
|
|
if (e.code() == error_code_end_of_stream) {
|
|
self->end.reset(self->messageVersion.version);
|
|
co_return;
|
|
}
|
|
throw e;
|
|
}
|
|
}
|
|
|
|
Future<Void> serverPeekGetMore(ServerPeekCursor* self, TaskPriority taskID) {
|
|
if (!self->interf || self->isExhausted()) {
|
|
return Never();
|
|
}
|
|
return serverPeekGetMoreImpl(self, taskID);
|
|
}
|
|
|
|
Future<Void> ServerPeekCursor::getMore(TaskPriority taskID) {
|
|
DebugLogTraceEvent("SPC_GetMore", randomID)
|
|
.detail("Tag", tag.toString())
|
|
.detail("HasMessage", hasMessage())
|
|
.detail("More", !more.isValid() || more.isReady())
|
|
.detail("Parallel", parallelGetMore)
|
|
.detail("MessageVersion", messageVersion.toString())
|
|
.detail("End", end.toString());
|
|
if (hasMessage() && !parallelGetMore) {
|
|
return Void();
|
|
}
|
|
if (!more.isValid() || more.isReady()) {
|
|
if (!hasMessage()) {
|
|
// A consumed reply is no longer useful. Release its arena before the next capped reply arrives so one
|
|
// ServerPeekCursor retains at most one reply window at a time.
|
|
results = TLogPeekReply();
|
|
results.maxKnownVersion = 0;
|
|
results.minKnownCommittedVersion = 0;
|
|
rd = ArenaReader(results.arena, results.messages, Unversioned());
|
|
}
|
|
if (usePeekStream &&
|
|
(tag.locality >= 0 || tag.locality == tagLocalityLogRouter || tag.locality == tagLocalityRemoteLog)) {
|
|
more = serverPeekStreamGetMore(this, taskID);
|
|
} else if (parallelGetMore || onlySpilled || !futureResults.empty()) {
|
|
more = serverPeekParallelGetMore(this, taskID);
|
|
} else {
|
|
more = serverPeekGetMore(this, taskID);
|
|
}
|
|
}
|
|
return more;
|
|
}
|
|
|
|
Future<Void> serverPeekOnFailed(ServerPeekCursor const* self) {
|
|
while (true) {
|
|
Future<Void> peekMessagesFailed = Never();
|
|
Future<Void> peekStreamMessagesFailed = Never();
|
|
if (self->interf->get().present()) {
|
|
peekMessagesFailed = IFailureMonitor::failureMonitor().onStateEqual(
|
|
self->interf->get().interf().peekMessages.getEndpoint(), FailureStatus());
|
|
peekStreamMessagesFailed = IFailureMonitor::failureMonitor().onStateEqual(
|
|
self->interf->get().interf().peekStreamMessages.getEndpoint(), FailureStatus());
|
|
}
|
|
auto res = co_await race(peekMessagesFailed, peekStreamMessagesFailed, self->interf->onChange());
|
|
if (res.index() == 0 || res.index() == 1) {
|
|
co_return;
|
|
} else if (res.index() != 2) {
|
|
UNREACHABLE();
|
|
}
|
|
}
|
|
}
|
|
|
|
Future<Void> ServerPeekCursor::onFailed() const {
|
|
return serverPeekOnFailed(this);
|
|
}
|
|
|
|
bool isAvailable(Reference<AsyncVar<OptionalInterface<TLogInterface>>> const& interf) {
|
|
if (!interf->get().present()) {
|
|
return false;
|
|
}
|
|
return IFailureMonitor::failureMonitor()
|
|
.getState(interf->get().interf().peekMessages.getEndpoint())
|
|
.isAvailable() &&
|
|
IFailureMonitor::failureMonitor()
|
|
.getState(interf->get().interf().peekStreamMessages.getEndpoint())
|
|
.isAvailable();
|
|
}
|
|
|
|
bool ServerPeekCursor::isActive() const {
|
|
if (isExhausted()) {
|
|
return false;
|
|
}
|
|
return isAvailable(interf);
|
|
}
|
|
|
|
bool ServerPeekCursor::isExhausted() const {
|
|
return messageVersion >= end;
|
|
}
|
|
|
|
const LogMessageVersion& ServerPeekCursor::version() const {
|
|
return messageVersion;
|
|
} // Call only after nextMessage(). The sequence of the current message, or results.end if nextMessage() has returned
|
|
// false.
|
|
|
|
Version ServerPeekCursor::getMinKnownCommittedVersion() const {
|
|
return results.minKnownCommittedVersion;
|
|
}
|
|
|
|
int64_t ServerPeekCursor::getMaxRetainedReplyCount() const {
|
|
return parallelGetMore || onlySpilled || !futureResults.empty()
|
|
? std::max<int64_t>(1, SERVER_KNOBS->PARALLEL_GET_MORE_REQUESTS + 1)
|
|
: 1;
|
|
}
|
|
|
|
void ServerPeekCursor::setReplyByteLimit(int limitBytes) {
|
|
ASSERT_GE(limitBytes, 0);
|
|
// A request that has already been issued cannot be retroactively capped.
|
|
ASSERT(!more.isValid());
|
|
ASSERT(futureResults.empty());
|
|
ASSERT(!peekReplyStream.present());
|
|
replyByteLimit = limitBytes;
|
|
}
|
|
|
|
Optional<UID> ServerPeekCursor::getPrimaryPeekLocation() const {
|
|
if (interf && interf->get().present()) {
|
|
return interf->get().id();
|
|
}
|
|
return Optional<UID>();
|
|
}
|
|
|
|
Optional<UID> ServerPeekCursor::getCurrentPeekLocation() const {
|
|
return ServerPeekCursor::getPrimaryPeekLocation();
|
|
}
|
|
|
|
Version ServerPeekCursor::popped() const {
|
|
return poppedVersion;
|
|
}
|
|
|
|
static void resetBestServerIfNotAvailable(
|
|
std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> const& logServers,
|
|
int& bestServer,
|
|
Version end) {
|
|
ASSERT(SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST);
|
|
if (bestServer >= 0 && end != std::numeric_limits<Version>::max()) {
|
|
if (!isAvailable(logServers[bestServer])) {
|
|
bestServer = -1;
|
|
}
|
|
}
|
|
}
|
|
|
|
/*
|
|
* Version vector/unicast related: This function is used by both set and merge peek cursors in order
|
|
* to decide when a tLog can return an empty version range. At a high level, the logic is as follows:
|
|
* - the tLogs of old epochs can return an empty version range
|
|
* - if "bestServer" is set then only the best server (of the "bestSet", if "bestSet" is initialized)
|
|
can return an empty version range
|
|
* - if "bestServer" is not set then only the tLogs that are known to have been locked can return an
|
|
* empty version range
|
|
* - if "bestSet" is set to a negative value then do not return an empty version range, irrespective
|
|
* of what other parameters are set to (we don't understand much about this scenario so trying to
|
|
* be safe in this case)
|
|
*/
|
|
static bool canReturnEmptyVersionRange(
|
|
int bestServer,
|
|
int currentServer,
|
|
Version end,
|
|
Optional<std::vector<uint16_t>> knownLockedTLogIndices = Optional<std::vector<uint16_t>>(),
|
|
Optional<int> bestSet = Optional<int>(),
|
|
Optional<int> currentSet = Optional<int>()) {
|
|
ASSERT(SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST && end != std::numeric_limits<Version>::max());
|
|
if (bestSet.present() && bestSet.get() < 0) {
|
|
// "BestSet" is set to a negative value - we don't know how/when this can happen.
|
|
// Be safe and return false.
|
|
return false;
|
|
}
|
|
if (!knownLockedTLogIndices.present()) {
|
|
// Servers from old epochs can return an empty version range.
|
|
return true;
|
|
}
|
|
bool foundServer =
|
|
std::binary_search(knownLockedTLogIndices.get().begin(), knownLockedTLogIndices.get().end(), currentServer);
|
|
if (bestServer >= 0) {
|
|
ASSERT_WE_THINK(!bestSet.present() || currentSet.present());
|
|
if ((!bestSet.present() || bestSet.get() == currentSet.get()) && currentServer == bestServer) {
|
|
ASSERT(foundServer);
|
|
// Best server is set - only the best server (that is known to have been locked) can return
|
|
// an empty version range.
|
|
return true;
|
|
}
|
|
} else if (foundServer) {
|
|
// Best server is not set - only servers that are known to have been locked can return
|
|
// an empty version range.
|
|
return true;
|
|
}
|
|
return false;
|
|
}
|
|
|
|
template <class SelectMessageVersion, class AdvanceAllCursors>
|
|
static bool updateMergedPeekMessage(std::vector<Reference<ServerPeekCursor>>& candidateCursors,
|
|
Optional<LogMessageVersion> const& nextVersion,
|
|
std::vector<std::pair<LogMessageVersion, int>>& sortedVersions,
|
|
LogMessageVersion& messageVersion,
|
|
int& currentCursor,
|
|
SelectMessageVersion selectMessageVersion,
|
|
AdvanceAllCursors advanceAllCursors) {
|
|
while (true) {
|
|
sortedVersions.clear();
|
|
for (int i = 0; i < candidateCursors.size(); i++) {
|
|
auto& candidateCursor = candidateCursors[i];
|
|
if (nextVersion.present()) {
|
|
candidateCursor->advanceTo(nextVersion.get());
|
|
}
|
|
sortedVersions.emplace_back(candidateCursor->version(), i);
|
|
}
|
|
|
|
selectMessageVersion(sortedVersions, messageVersion);
|
|
if (!advanceAllCursors(messageVersion)) {
|
|
break;
|
|
}
|
|
}
|
|
|
|
for (int i = 0; i < candidateCursors.size(); i++) {
|
|
auto& candidateCursor = candidateCursors[i];
|
|
ASSERT_WE_THINK(!candidateCursor->hasMessage() || candidateCursor->version() >= messageVersion);
|
|
if (candidateCursor->version() == messageVersion && candidateCursor->hasMessage()) {
|
|
currentCursor = i;
|
|
return true;
|
|
}
|
|
}
|
|
return false;
|
|
}
|
|
|
|
MergedPeekCursor::MergedPeekCursor(std::vector<Reference<ServerPeekCursor>> const& serverCursors, Version begin)
|
|
: serverCursors(serverCursors), tag(invalidTag), bestServer(-1), currentCursor(0), readQuorum(serverCursors.size()),
|
|
messageVersion(begin), hasNextMessage(false), randomID(deterministicRandom()->randomUniqueID()),
|
|
tLogReplicationFactor(0) {
|
|
sortedVersions.resize(serverCursors.size());
|
|
}
|
|
|
|
MergedPeekCursor::MergedPeekCursor(std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> const& logServers,
|
|
int bestServerLogId,
|
|
int readQuorum,
|
|
Tag tag,
|
|
Version begin,
|
|
Version end,
|
|
bool parallelGetMore,
|
|
std::vector<LocalityData> const& tLogLocalities,
|
|
Reference<IReplicationPolicy> const tLogPolicy,
|
|
int tLogReplicationFactor,
|
|
const Optional<std::vector<uint16_t>>& knownLockedTLogIds)
|
|
: tag(tag), bestServer(bestServerLogId), currentCursor(0), readQuorum(readQuorum), messageVersion(begin),
|
|
hasNextMessage(false), randomID(deterministicRandom()->randomUniqueID()),
|
|
tLogReplicationFactor(tLogReplicationFactor) {
|
|
if (tLogPolicy) {
|
|
logSet = makeReference<LogSet>();
|
|
logSet->tLogPolicy = tLogPolicy;
|
|
logSet->tLogLocalities = tLogLocalities;
|
|
filterLocalityDataForPolicy(logSet->tLogPolicy, &logSet->tLogLocalities);
|
|
logSet->updateLocalitySet(logSet->tLogLocalities);
|
|
}
|
|
|
|
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
|
|
resetBestServerIfNotAvailable(logServers, bestServer, end);
|
|
}
|
|
|
|
for (int i = 0; i < logServers.size(); i++) {
|
|
bool returnEmptyIfStopped =
|
|
((SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST && end != std::numeric_limits<Version>::max())
|
|
? canReturnEmptyVersionRange(bestServer, i /*currentServer*/, end, knownLockedTLogIds)
|
|
: false);
|
|
auto cursor = makeReference<ServerPeekCursor>(
|
|
logServers[i], tag, begin, end, bestServer >= 0, parallelGetMore, returnEmptyIfStopped);
|
|
//TraceEvent("MPC_Starting", randomID).detail("Cursor", cursor->randomID).detail("End", end);
|
|
serverCursors.push_back(cursor);
|
|
}
|
|
sortedVersions.resize(serverCursors.size());
|
|
}
|
|
|
|
MergedPeekCursor::MergedPeekCursor(std::vector<Reference<ServerPeekCursor>> const& serverCursors,
|
|
LogMessageVersion const& messageVersion,
|
|
int bestServer,
|
|
int readQuorum,
|
|
Optional<LogMessageVersion> nextVersion,
|
|
Reference<LogSet> logSet,
|
|
int tLogReplicationFactor)
|
|
: logSet(logSet), serverCursors(serverCursors), bestServer(bestServer), currentCursor(0), readQuorum(readQuorum),
|
|
nextVersion(nextVersion), messageVersion(messageVersion), hasNextMessage(false),
|
|
randomID(deterministicRandom()->randomUniqueID()), tLogReplicationFactor(tLogReplicationFactor) {
|
|
sortedVersions.resize(serverCursors.size());
|
|
calcHasMessage();
|
|
}
|
|
|
|
Reference<IReplayPeekCursor> MergedPeekCursor::cloneNoMore() {
|
|
std::vector<Reference<ServerPeekCursor>> cursors;
|
|
for (const auto& it : serverCursors) {
|
|
cursors.push_back(it->cloneServerNoMore());
|
|
}
|
|
return makeReference<MergedPeekCursor>(
|
|
cursors, messageVersion, bestServer, readQuorum, nextVersion, logSet, tLogReplicationFactor);
|
|
}
|
|
|
|
void MergedPeekCursor::setProtocolVersion(ProtocolVersion version) {
|
|
for (const auto& it : serverCursors) {
|
|
if (it->hasMessage()) {
|
|
it->setProtocolVersion(version);
|
|
}
|
|
}
|
|
}
|
|
|
|
Arena& MergedPeekCursor::arena() {
|
|
return serverCursors[currentCursor]->arena();
|
|
}
|
|
|
|
ArenaReader* MergedPeekCursor::reader() {
|
|
return serverCursors[currentCursor]->reader();
|
|
}
|
|
|
|
void MergedPeekCursor::calcHasMessage() {
|
|
if (bestServer >= 0) {
|
|
if (nextVersion.present()) {
|
|
serverCursors[bestServer]->advanceTo(nextVersion.get());
|
|
}
|
|
if (serverCursors[bestServer]->hasMessage()) {
|
|
messageVersion = serverCursors[bestServer]->version();
|
|
currentCursor = bestServer;
|
|
hasNextMessage = true;
|
|
|
|
for (auto& c : serverCursors) {
|
|
c->advanceTo(messageVersion);
|
|
}
|
|
|
|
return;
|
|
}
|
|
|
|
auto bestVersion = serverCursors[bestServer]->version();
|
|
for (auto& c : serverCursors) {
|
|
c->advanceTo(bestVersion);
|
|
}
|
|
}
|
|
|
|
hasNextMessage = false;
|
|
updateMessage(false);
|
|
|
|
if (!hasNextMessage && logSet) {
|
|
updateMessage(true);
|
|
}
|
|
}
|
|
|
|
void MergedPeekCursor::updateMessage(bool usePolicy) {
|
|
auto selectMessageVersion = [&](auto& versions, LogMessageVersion& selectedVersion) {
|
|
if (usePolicy) {
|
|
ASSERT(logSet->tLogPolicy);
|
|
std::sort(versions.begin(), versions.end());
|
|
|
|
locations.clear();
|
|
for (auto sortedVersion : versions) {
|
|
locations.push_back(logSet->logEntryArray[sortedVersion.second]);
|
|
if (locations.size() >= tLogReplicationFactor && logSet->satisfiesPolicy(locations)) {
|
|
selectedVersion = sortedVersion.first;
|
|
break;
|
|
}
|
|
}
|
|
} else {
|
|
std::nth_element(versions.begin(), versions.end() - readQuorum, versions.end());
|
|
selectedVersion = versions[versions.size() - readQuorum].first;
|
|
}
|
|
};
|
|
auto advanceAllCursors = [&](LogMessageVersion const& selectedVersion) {
|
|
bool advancedPast = false;
|
|
for (auto& cursor : serverCursors) {
|
|
auto start = cursor->version();
|
|
cursor->advanceTo(selectedVersion);
|
|
if (start <= selectedVersion && selectedVersion < cursor->version()) {
|
|
advancedPast = true;
|
|
CODE_PROBE(true, "Merge peek cursor advanced past desired sequence", probe::decoration::rare);
|
|
}
|
|
}
|
|
return advancedPast;
|
|
};
|
|
|
|
hasNextMessage = updateMergedPeekMessage(serverCursors,
|
|
nextVersion,
|
|
sortedVersions,
|
|
messageVersion,
|
|
currentCursor,
|
|
selectMessageVersion,
|
|
advanceAllCursors);
|
|
}
|
|
|
|
bool MergedPeekCursor::hasMessage() const {
|
|
return hasNextMessage;
|
|
}
|
|
|
|
void MergedPeekCursor::nextMessage() {
|
|
nextVersion = version();
|
|
nextVersion.get().sub++;
|
|
serverCursors[currentCursor]->nextMessage();
|
|
calcHasMessage();
|
|
ASSERT(hasMessage() || !version().sub);
|
|
}
|
|
|
|
StringRef MergedPeekCursor::getMessage() {
|
|
return serverCursors[currentCursor]->getMessage();
|
|
}
|
|
|
|
StringRef MergedPeekCursor::getMessageWithTags() {
|
|
return serverCursors[currentCursor]->getMessageWithTags();
|
|
}
|
|
|
|
VectorRef<Tag> MergedPeekCursor::getTags() const {
|
|
return serverCursors[currentCursor]->getTags();
|
|
}
|
|
|
|
void MergedPeekCursor::advanceTo(LogMessageVersion n) {
|
|
bool canChange = false;
|
|
for (auto& c : serverCursors) {
|
|
if (c->version() < n) {
|
|
canChange = true;
|
|
c->advanceTo(n);
|
|
}
|
|
}
|
|
if (canChange) {
|
|
calcHasMessage();
|
|
}
|
|
}
|
|
|
|
Future<Void> mergedPeekGetMore(MergedPeekCursor* self, LogMessageVersion startVersion, TaskPriority taskID) {
|
|
while (true) {
|
|
//TraceEvent("MPC_GetMoreA", self->randomID).detail("Start", startVersion.toString());
|
|
if (self->bestServer >= 0 && self->serverCursors[self->bestServer]->isActive()) {
|
|
ASSERT(!self->serverCursors[self->bestServer]->hasMessage());
|
|
co_await (self->serverCursors[self->bestServer]->getMore(taskID) ||
|
|
self->serverCursors[self->bestServer]->onFailed());
|
|
} else {
|
|
std::vector<Future<Void>> q;
|
|
for (auto& c : self->serverCursors) {
|
|
if (!c->hasMessage()) {
|
|
q.push_back(c->getMore(taskID));
|
|
}
|
|
}
|
|
co_await quorum(q, 1);
|
|
}
|
|
self->calcHasMessage();
|
|
//TraceEvent("MPC_GetMoreB", self->randomID).detail("HasMessage", self->hasMessage()).detail("Start", startVersion.toString()).detail("Seq", self->version().toString());
|
|
if (self->hasMessage() || self->version() > startVersion) {
|
|
co_return;
|
|
}
|
|
}
|
|
}
|
|
|
|
Future<Void> MergedPeekCursor::getMore(TaskPriority taskID) {
|
|
if (more.isValid() && !more.isReady()) {
|
|
return more;
|
|
}
|
|
|
|
if (serverCursors.empty()) {
|
|
return Never();
|
|
}
|
|
|
|
auto startVersion = version();
|
|
calcHasMessage();
|
|
if (hasMessage()) {
|
|
return Void();
|
|
}
|
|
if (nextVersion.present()) {
|
|
advanceTo(nextVersion.get());
|
|
}
|
|
ASSERT(!hasMessage());
|
|
if (version() > startVersion) {
|
|
return Void();
|
|
}
|
|
|
|
more = mergedPeekGetMore(this, startVersion, taskID);
|
|
return more;
|
|
}
|
|
|
|
bool MergedPeekCursor::isExhausted() const {
|
|
return serverCursors[currentCursor]->isExhausted();
|
|
}
|
|
|
|
const LogMessageVersion& MergedPeekCursor::version() const {
|
|
return messageVersion;
|
|
}
|
|
|
|
Version MergedPeekCursor::getMinKnownCommittedVersion() const {
|
|
return serverCursors[currentCursor]->getMinKnownCommittedVersion();
|
|
}
|
|
|
|
Version MergedPeekCursor::getMaxKnownVersion() const {
|
|
Version maxKnownVersion = 0;
|
|
for (const auto& cursor : serverCursors) {
|
|
maxKnownVersion = std::max(maxKnownVersion, cursor->getMaxKnownVersion());
|
|
}
|
|
return maxKnownVersion;
|
|
}
|
|
|
|
int64_t MergedPeekCursor::getMaxRetainedReplyCount() const {
|
|
int64_t count = 0;
|
|
for (const auto& cursor : serverCursors) {
|
|
count += cursor->getMaxRetainedReplyCount();
|
|
}
|
|
return std::max<int64_t>(1, count);
|
|
}
|
|
|
|
void MergedPeekCursor::setReplyByteLimit(int limitBytes) {
|
|
for (const auto& cursor : serverCursors) {
|
|
cursor->setReplyByteLimit(limitBytes);
|
|
}
|
|
}
|
|
|
|
Optional<UID> MergedPeekCursor::getPrimaryPeekLocation() const {
|
|
if (bestServer >= 0) {
|
|
return serverCursors[bestServer]->getPrimaryPeekLocation();
|
|
}
|
|
return Optional<UID>();
|
|
}
|
|
|
|
Optional<UID> MergedPeekCursor::getCurrentPeekLocation() const {
|
|
if (currentCursor >= 0) {
|
|
return serverCursors[currentCursor]->getPrimaryPeekLocation();
|
|
}
|
|
return Optional<UID>();
|
|
}
|
|
|
|
Version MergedPeekCursor::popped() const {
|
|
Version poppedVersion = 0;
|
|
for (auto& c : serverCursors) {
|
|
poppedVersion = std::max(poppedVersion, c->popped());
|
|
}
|
|
return poppedVersion;
|
|
}
|
|
|
|
SetPeekCursor::SetPeekCursor(std::vector<Reference<LogSet>> const& logSets,
|
|
int bestSet,
|
|
int bestServerLogId,
|
|
Tag tag,
|
|
Version begin,
|
|
Version end,
|
|
bool parallelGetMore,
|
|
const Optional<std::vector<uint16_t>>& knownLockedTLogIds)
|
|
: logSets(logSets), tag(tag), bestSet(bestSet), bestServer(bestServerLogId), currentSet(bestSet), currentCursor(0),
|
|
messageVersion(begin), hasNextMessage(false), useBestSet(true), randomID(deterministicRandom()->randomUniqueID()),
|
|
end(end) {
|
|
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST && bestSet >= 0) {
|
|
resetBestServerIfNotAvailable(logSets[bestSet]->logServers, bestServer, end);
|
|
}
|
|
// Live CDC consumers wait for future mutations; finite recovery and tag-history reads must reach their end.
|
|
const bool tailingCDC = tag.locality == tagLocalityCDC && end == std::numeric_limits<Version>::max();
|
|
CODE_PROBE(tag.locality == tagLocalityCDC && !tailingCDC, "CDC finite-range peek returns without blocking");
|
|
serverCursors.resize(logSets.size());
|
|
int maxServers = 0;
|
|
for (int i = 0; i < logSets.size(); i++) {
|
|
for (int j = 0; j < logSets[i]->logServers.size(); j++) {
|
|
bool returnEmptyIfStopped =
|
|
((SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST && end != std::numeric_limits<Version>::max())
|
|
? canReturnEmptyVersionRange(
|
|
bestServer, j /*currentServer*/, end, knownLockedTLogIds, bestSet, i /* currentSet */)
|
|
: false);
|
|
auto cursor = makeReference<ServerPeekCursor>(
|
|
logSets[i]->logServers[j], tag, begin, end, !tailingCDC, parallelGetMore, returnEmptyIfStopped);
|
|
serverCursors[i].push_back(cursor);
|
|
}
|
|
maxServers = std::max<int>(maxServers, serverCursors[i].size());
|
|
}
|
|
sortedVersions.resize(maxServers);
|
|
}
|
|
|
|
SetPeekCursor::SetPeekCursor(std::vector<Reference<LogSet>> const& logSets,
|
|
std::vector<std::vector<Reference<ServerPeekCursor>>> const& serverCursors,
|
|
LogMessageVersion const& messageVersion,
|
|
int bestSet,
|
|
int bestServer,
|
|
Optional<LogMessageVersion> nextVersion,
|
|
bool useBestSet)
|
|
: logSets(logSets), serverCursors(serverCursors), bestSet(bestSet), bestServer(bestServer), currentSet(bestSet),
|
|
currentCursor(0), nextVersion(nextVersion), messageVersion(messageVersion), hasNextMessage(false),
|
|
useBestSet(useBestSet), randomID(deterministicRandom()->randomUniqueID()) {
|
|
int maxServers = 0;
|
|
for (int i = 0; i < logSets.size(); i++) {
|
|
maxServers = std::max<int>(maxServers, serverCursors[i].size());
|
|
}
|
|
sortedVersions.resize(maxServers);
|
|
calcHasMessage();
|
|
}
|
|
|
|
Reference<IReplayPeekCursor> SetPeekCursor::cloneNoMore() {
|
|
std::vector<std::vector<Reference<ServerPeekCursor>>> cursors;
|
|
cursors.resize(logSets.size());
|
|
for (int i = 0; i < logSets.size(); i++) {
|
|
for (int j = 0; j < logSets[i]->logServers.size(); j++) {
|
|
cursors[i].push_back(serverCursors[i][j]->cloneServerNoMore());
|
|
}
|
|
}
|
|
return makeReference<SetPeekCursor>(logSets, cursors, messageVersion, bestSet, bestServer, nextVersion, useBestSet);
|
|
}
|
|
|
|
void SetPeekCursor::setProtocolVersion(ProtocolVersion version) {
|
|
for (auto& cursors : serverCursors) {
|
|
for (auto& it : cursors) {
|
|
if (it->hasMessage()) {
|
|
it->setProtocolVersion(version);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
Arena& SetPeekCursor::arena() {
|
|
return serverCursors[currentSet][currentCursor]->arena();
|
|
}
|
|
|
|
ArenaReader* SetPeekCursor::reader() {
|
|
return serverCursors[currentSet][currentCursor]->reader();
|
|
}
|
|
|
|
void SetPeekCursor::calcHasMessage() {
|
|
if (bestSet >= 0 && bestServer >= 0) {
|
|
if (nextVersion.present()) {
|
|
//TraceEvent("LPC_CalcNext").detail("Ver", messageVersion.toString()).detail("Tag", tag.toString()).detail("HasNextMessage", hasNextMessage).detail("NextVersion", nextVersion.get().toString());
|
|
serverCursors[bestSet][bestServer]->advanceTo(nextVersion.get());
|
|
}
|
|
if (serverCursors[bestSet][bestServer]->hasMessage()) {
|
|
messageVersion = serverCursors[bestSet][bestServer]->version();
|
|
currentSet = bestSet;
|
|
currentCursor = bestServer;
|
|
hasNextMessage = true;
|
|
|
|
//TraceEvent("LPC_Calc1").detail("Ver", messageVersion.toString()).detail("Tag", tag.toString()).detail("HasNextMessage", hasNextMessage);
|
|
|
|
for (auto& cursors : serverCursors) {
|
|
for (auto& c : cursors) {
|
|
c->advanceTo(messageVersion);
|
|
}
|
|
}
|
|
|
|
return;
|
|
}
|
|
|
|
auto bestVersion = serverCursors[bestSet][bestServer]->version();
|
|
for (auto& cursors : serverCursors) {
|
|
for (auto& c : cursors) {
|
|
c->advanceTo(bestVersion);
|
|
}
|
|
}
|
|
}
|
|
|
|
hasNextMessage = false;
|
|
if (useBestSet) {
|
|
updateMessage(bestSet, false); // Use Quorum logic
|
|
|
|
//TraceEvent("LPC_Calc2").detail("Ver", messageVersion.toString()).detail("Tag", tag.toString()).detail("HasNextMessage", hasNextMessage);
|
|
if (!hasNextMessage) {
|
|
updateMessage(bestSet, true);
|
|
//TraceEvent("LPC_Calc3").detail("Ver", messageVersion.toString()).detail("Tag", tag.toString()).detail("HasNextMessage", hasNextMessage);
|
|
}
|
|
} else {
|
|
for (int i = 0; i < logSets.size() && !hasNextMessage; i++) {
|
|
if (i != bestSet) {
|
|
updateMessage(i, false); // Use Quorum logic
|
|
}
|
|
}
|
|
//TraceEvent("LPC_Calc4").detail("Ver", messageVersion.toString()).detail("Tag", tag.toString()).detail("HasNextMessage", hasNextMessage);
|
|
for (int i = 0; i < logSets.size() && !hasNextMessage; i++) {
|
|
if (i != bestSet) {
|
|
updateMessage(i, true);
|
|
}
|
|
}
|
|
//TraceEvent("LPC_Calc5").detail("Ver", messageVersion.toString()).detail("Tag", tag.toString()).detail("HasNextMessage", hasNextMessage);
|
|
}
|
|
}
|
|
|
|
void SetPeekCursor::updateMessage(int logIdx, bool usePolicy) {
|
|
auto selectMessageVersion = [&](auto& versions, LogMessageVersion& selectedVersion) {
|
|
if (usePolicy) {
|
|
std::sort(versions.begin(), versions.end());
|
|
locations.clear();
|
|
for (auto sortedVersion : versions) {
|
|
locations.push_back(logSets[logIdx]->logEntryArray[sortedVersion.second]);
|
|
if (locations.size() >= logSets[logIdx]->tLogReplicationFactor &&
|
|
logSets[logIdx]->satisfiesPolicy(locations)) {
|
|
selectedVersion = sortedVersion.first;
|
|
break;
|
|
}
|
|
}
|
|
} else {
|
|
//(int)oldLogData[i].logServers.size() + 1 - oldLogData[i].tLogReplicationFactor
|
|
auto quorum = logSets[logIdx]->logServers.size() + 1 - logSets[logIdx]->tLogReplicationFactor;
|
|
std::nth_element(versions.begin(), versions.end() - quorum, versions.end());
|
|
selectedVersion = versions[versions.size() - quorum].first;
|
|
}
|
|
};
|
|
auto advanceAllCursors = [&](LogMessageVersion const& selectedVersion) {
|
|
bool advancedPast = false;
|
|
for (auto& cursors : serverCursors) {
|
|
for (auto& cursor : cursors) {
|
|
auto start = cursor->version();
|
|
cursor->advanceTo(selectedVersion);
|
|
if (start <= selectedVersion && selectedVersion < cursor->version()) {
|
|
advancedPast = true;
|
|
CODE_PROBE(true, "Set peek cursor with logIdx advanced past desired sequence");
|
|
}
|
|
}
|
|
}
|
|
return advancedPast;
|
|
};
|
|
|
|
hasNextMessage = updateMergedPeekMessage(serverCursors[logIdx],
|
|
nextVersion,
|
|
sortedVersions,
|
|
messageVersion,
|
|
currentCursor,
|
|
selectMessageVersion,
|
|
advanceAllCursors);
|
|
if (hasNextMessage) {
|
|
currentSet = logIdx;
|
|
}
|
|
}
|
|
|
|
bool SetPeekCursor::hasMessage() const {
|
|
return hasNextMessage;
|
|
}
|
|
|
|
void SetPeekCursor::nextMessage() {
|
|
nextVersion = version();
|
|
nextVersion.get().sub++;
|
|
serverCursors[currentSet][currentCursor]->nextMessage();
|
|
calcHasMessage();
|
|
ASSERT(hasMessage() || !version().sub);
|
|
}
|
|
|
|
StringRef SetPeekCursor::getMessage() {
|
|
return serverCursors[currentSet][currentCursor]->getMessage();
|
|
}
|
|
|
|
StringRef SetPeekCursor::getMessageWithTags() {
|
|
return serverCursors[currentSet][currentCursor]->getMessageWithTags();
|
|
}
|
|
|
|
VectorRef<Tag> SetPeekCursor::getTags() const {
|
|
return serverCursors[currentSet][currentCursor]->getTags();
|
|
}
|
|
|
|
void SetPeekCursor::advanceTo(LogMessageVersion n) {
|
|
bool canChange = false;
|
|
for (auto& cursors : serverCursors) {
|
|
for (auto& c : cursors) {
|
|
if (c->version() < n) {
|
|
canChange = true;
|
|
c->advanceTo(n);
|
|
}
|
|
}
|
|
}
|
|
if (canChange) {
|
|
calcHasMessage();
|
|
}
|
|
}
|
|
|
|
Future<Void> setPeekGetMore(SetPeekCursor* self, LogMessageVersion startVersion, TaskPriority taskID) {
|
|
while (true) {
|
|
//TraceEvent("LPC_GetMore1", self->randomID).detail("Start", startVersion.toString()).detail("Tag", self->tag.toString());
|
|
if (self->bestServer >= 0 && self->bestSet >= 0 &&
|
|
self->serverCursors[self->bestSet][self->bestServer]->isActive()) {
|
|
ASSERT(!self->serverCursors[self->bestSet][self->bestServer]->hasMessage());
|
|
//TraceEvent("LPC_GetMore2", self->randomID).detail("Start", startVersion.toString()).detail("Tag", self->tag.toString());
|
|
co_await (self->serverCursors[self->bestSet][self->bestServer]->getMore(taskID) ||
|
|
self->serverCursors[self->bestSet][self->bestServer]->onFailed());
|
|
self->useBestSet = true;
|
|
} else {
|
|
// FIXME: if best set is exhausted, do not peek remote servers
|
|
bool bestSetValid = self->bestSet >= 0;
|
|
if (bestSetValid) {
|
|
self->locations.clear();
|
|
for (int i = 0; i < self->serverCursors[self->bestSet].size(); i++) {
|
|
if (!self->serverCursors[self->bestSet][i]->isActive() &&
|
|
self->serverCursors[self->bestSet][i]->version() <= self->messageVersion) {
|
|
self->locations.push_back(self->logSets[self->bestSet]->logEntryArray[i]);
|
|
}
|
|
}
|
|
bestSetValid = self->locations.size() < self->logSets[self->bestSet]->tLogReplicationFactor ||
|
|
!self->logSets[self->bestSet]->satisfiesPolicy(self->locations);
|
|
}
|
|
if (bestSetValid || self->logSets.size() == 1) {
|
|
if (!self->useBestSet) {
|
|
self->useBestSet = true;
|
|
self->calcHasMessage();
|
|
if (self->hasMessage() || self->version() > startVersion) {
|
|
co_return;
|
|
}
|
|
}
|
|
|
|
//TraceEvent("LPC_GetMore3", self->randomID).detail("Start", startVersion.toString()).detail("Tag", self->tag.toString()).detail("BestSetSize", self->serverCursors[self->bestSet].size());
|
|
std::vector<Future<Void>> q;
|
|
for (auto& c : self->serverCursors[self->bestSet]) {
|
|
if (!c->hasMessage()) {
|
|
q.push_back(c->getMore(taskID));
|
|
if (c->isActive()) {
|
|
q.push_back(c->onFailed());
|
|
}
|
|
}
|
|
}
|
|
co_await quorum(q, 1);
|
|
} else {
|
|
// FIXME: this will peeking way too many cursors when satellites exist, and does not need to peek
|
|
// bestSet cursors since we cannot get anymore data from them
|
|
std::vector<Future<Void>> q;
|
|
//TraceEvent("LPC_GetMore4", self->randomID).detail("Start", startVersion.toString()).detail("Tag", self->tag.toString());
|
|
for (auto& cursors : self->serverCursors) {
|
|
for (auto& c : cursors) {
|
|
if (!c->hasMessage()) {
|
|
q.push_back(c->getMore(taskID));
|
|
}
|
|
}
|
|
}
|
|
co_await quorum(q, 1);
|
|
self->useBestSet = false;
|
|
}
|
|
}
|
|
self->calcHasMessage();
|
|
//TraceEvent("LPC_GetMoreB", self->randomID).detail("HasMessage", self->hasMessage()).detail("Start", startVersion.toString()).detail("Seq", self->version().toString());
|
|
if (self->hasMessage() || self->version() > startVersion) {
|
|
co_return;
|
|
}
|
|
}
|
|
}
|
|
|
|
Future<Void> SetPeekCursor::getMore(TaskPriority taskID) {
|
|
if (more.isValid() && !more.isReady()) {
|
|
return more;
|
|
}
|
|
|
|
auto startVersion = version();
|
|
calcHasMessage();
|
|
if (hasMessage()) {
|
|
return Void();
|
|
}
|
|
if (nextVersion.present()) {
|
|
advanceTo(nextVersion.get());
|
|
}
|
|
ASSERT(!hasMessage());
|
|
if (version() > startVersion) {
|
|
return Void();
|
|
}
|
|
|
|
more = setPeekGetMore(this, startVersion, taskID);
|
|
return more;
|
|
}
|
|
|
|
bool SetPeekCursor::isExhausted() const {
|
|
return serverCursors[currentSet][currentCursor]->isExhausted();
|
|
}
|
|
|
|
const LogMessageVersion& SetPeekCursor::version() const {
|
|
return messageVersion;
|
|
}
|
|
|
|
Version SetPeekCursor::getMinKnownCommittedVersion() const {
|
|
return serverCursors[currentSet][currentCursor]->getMinKnownCommittedVersion();
|
|
}
|
|
|
|
Version SetPeekCursor::getMaxKnownVersion() const {
|
|
Version maxKnownVersion = 0;
|
|
for (const auto& cursors : serverCursors) {
|
|
for (const auto& cursor : cursors) {
|
|
maxKnownVersion = std::max(maxKnownVersion, cursor->getMaxKnownVersion());
|
|
}
|
|
}
|
|
return maxKnownVersion;
|
|
}
|
|
|
|
int64_t SetPeekCursor::getMaxRetainedReplyCount() const {
|
|
int64_t count = 0;
|
|
for (const auto& cursors : serverCursors) {
|
|
for (const auto& cursor : cursors) {
|
|
count += cursor->getMaxRetainedReplyCount();
|
|
}
|
|
}
|
|
return std::max<int64_t>(1, count);
|
|
}
|
|
|
|
void SetPeekCursor::setReplyByteLimit(int limitBytes) {
|
|
for (const auto& cursors : serverCursors) {
|
|
for (const auto& cursor : cursors) {
|
|
cursor->setReplyByteLimit(limitBytes);
|
|
}
|
|
}
|
|
}
|
|
|
|
Optional<UID> SetPeekCursor::getPrimaryPeekLocation() const {
|
|
if (bestServer >= 0 && bestSet >= 0) {
|
|
return serverCursors[bestSet][bestServer]->getPrimaryPeekLocation();
|
|
}
|
|
return Optional<UID>();
|
|
}
|
|
|
|
Optional<UID> SetPeekCursor::getCurrentPeekLocation() const {
|
|
if (currentCursor >= 0 && currentSet >= 0) {
|
|
return serverCursors[currentSet][currentCursor]->getPrimaryPeekLocation();
|
|
}
|
|
return Optional<UID>();
|
|
}
|
|
|
|
Version SetPeekCursor::popped() const {
|
|
Version poppedVersion = 0;
|
|
for (auto& cursors : serverCursors) {
|
|
for (auto& c : cursors) {
|
|
poppedVersion = std::max(poppedVersion, c->popped());
|
|
}
|
|
}
|
|
return poppedVersion;
|
|
}
|
|
|
|
ReplayMultiCursor::ReplayMultiCursor(std::vector<Reference<IReplayPeekCursor>> cursors,
|
|
std::vector<LogMessageVersion> epochEnds,
|
|
bool prefetch)
|
|
: cursors(cursors), epochEnds(epochEnds), poppedVersion(0), prefetch(prefetch) {
|
|
if (prefetch) {
|
|
for (int i = 0; i < std::min<int>(cursors.size(), SERVER_KNOBS->MULTI_CURSOR_PRE_FETCH_LIMIT); i++) {
|
|
cursors[cursors.size() - i - 1]->getMore();
|
|
}
|
|
}
|
|
}
|
|
|
|
Reference<IReplayPeekCursor> ReplayMultiCursor::cloneNoMore() {
|
|
return cursors.back()->cloneNoMore();
|
|
}
|
|
|
|
void ReplayMultiCursor::setProtocolVersion(ProtocolVersion version) {
|
|
cursors.back()->setProtocolVersion(version);
|
|
}
|
|
|
|
Arena& ReplayMultiCursor::arena() {
|
|
return cursors.back()->arena();
|
|
}
|
|
|
|
ArenaReader* ReplayMultiCursor::reader() {
|
|
return cursors.back()->reader();
|
|
}
|
|
|
|
bool ReplayMultiCursor::hasMessage() const {
|
|
return cursors.back()->hasMessage();
|
|
}
|
|
|
|
void ReplayMultiCursor::nextMessage() {
|
|
cursors.back()->nextMessage();
|
|
}
|
|
|
|
StringRef ReplayMultiCursor::getMessage() {
|
|
return cursors.back()->getMessage();
|
|
}
|
|
|
|
StringRef ReplayMultiCursor::getMessageWithTags() {
|
|
return cursors.back()->getMessageWithTags();
|
|
}
|
|
|
|
VectorRef<Tag> ReplayMultiCursor::getTags() const {
|
|
return cursors.back()->getTags();
|
|
}
|
|
|
|
void ReplayMultiCursor::advanceTo(LogMessageVersion n) {
|
|
while (cursors.size() > 1 && n >= epochEnds.back()) {
|
|
poppedVersion = std::max(poppedVersion, cursors.back()->popped());
|
|
cursors.pop_back();
|
|
epochEnds.pop_back();
|
|
}
|
|
cursors.back()->advanceTo(n);
|
|
}
|
|
|
|
Future<Void> ReplayMultiCursor::getMore(TaskPriority taskID) {
|
|
LogMessageVersion startVersion = cursors.back()->version();
|
|
while (cursors.size() > 1 && cursors.back()->version() >= epochEnds.back()) {
|
|
poppedVersion = std::max(poppedVersion, cursors.back()->popped());
|
|
cursors.pop_back();
|
|
epochEnds.pop_back();
|
|
}
|
|
if (cursors.back()->version() > startVersion) {
|
|
return Void();
|
|
}
|
|
return cursors.back()->getMore(taskID);
|
|
}
|
|
|
|
bool ReplayMultiCursor::isExhausted() const {
|
|
return cursors.back()->isExhausted();
|
|
}
|
|
|
|
const LogMessageVersion& ReplayMultiCursor::version() const {
|
|
return cursors.back()->version();
|
|
}
|
|
|
|
Version ReplayMultiCursor::getMinKnownCommittedVersion() const {
|
|
const Version cursorCommittedVersion = cursors.back()->getMinKnownCommittedVersion();
|
|
if (cursors.size() == 1) {
|
|
return cursorCommittedVersion;
|
|
}
|
|
|
|
// A following generation starts one version after the completed generation's known committed frontier. A
|
|
// stopped TLog can report an older local frontier forever, so expose the generation boundary to readers that
|
|
// gate delivery on committed progress.
|
|
const Version completedGenerationCommittedVersion = epochEnds.back().version - 1;
|
|
CODE_PROBE(cursorCommittedVersion < completedGenerationCommittedVersion,
|
|
"Replay cursor advances the committed frontier through a completed generation");
|
|
return std::max(cursorCommittedVersion, completedGenerationCommittedVersion);
|
|
}
|
|
|
|
Version ReplayMultiCursor::getMaxKnownVersion() const {
|
|
Version maxKnownVersion = 0;
|
|
for (const auto& cursor : cursors) {
|
|
maxKnownVersion = std::max(maxKnownVersion, cursor->getMaxKnownVersion());
|
|
}
|
|
return maxKnownVersion;
|
|
}
|
|
|
|
int64_t ReplayMultiCursor::getMaxRetainedReplyCount() const {
|
|
int64_t count = prefetch ? 0 : 1;
|
|
for (const auto& cursor : cursors) {
|
|
if (prefetch) {
|
|
count += cursor->getMaxRetainedReplyCount();
|
|
} else {
|
|
count = std::max(count, cursor->getMaxRetainedReplyCount());
|
|
}
|
|
}
|
|
return count;
|
|
}
|
|
|
|
void ReplayMultiCursor::setReplyByteLimit(int limitBytes) {
|
|
for (const auto& cursor : cursors) {
|
|
cursor->setReplyByteLimit(limitBytes);
|
|
}
|
|
}
|
|
|
|
Optional<UID> ReplayMultiCursor::getPrimaryPeekLocation() const {
|
|
return cursors.back()->getPrimaryPeekLocation();
|
|
}
|
|
|
|
Optional<UID> ReplayMultiCursor::getCurrentPeekLocation() const {
|
|
return cursors.back()->getCurrentPeekLocation();
|
|
}
|
|
|
|
Version ReplayMultiCursor::popped() const {
|
|
return std::max(poppedVersion, cursors.back()->popped());
|
|
}
|
|
|
|
MultiCursor::MultiCursor(std::vector<Reference<IPeekCursor>> cursors, std::vector<LogMessageVersion> epochEnds)
|
|
: cursors(cursors), epochEnds(epochEnds), poppedVersion(0) {
|
|
for (int i = 0; i < std::min<int>(cursors.size(), SERVER_KNOBS->MULTI_CURSOR_PRE_FETCH_LIMIT); i++) {
|
|
cursors[cursors.size() - i - 1]->getMore();
|
|
}
|
|
}
|
|
|
|
void MultiCursor::setProtocolVersion(ProtocolVersion version) {
|
|
cursors.back()->setProtocolVersion(version);
|
|
}
|
|
|
|
Arena& MultiCursor::arena() {
|
|
return cursors.back()->arena();
|
|
}
|
|
|
|
ArenaReader* MultiCursor::reader() {
|
|
return cursors.back()->reader();
|
|
}
|
|
|
|
bool MultiCursor::hasMessage() const {
|
|
return cursors.back()->hasMessage();
|
|
}
|
|
|
|
void MultiCursor::nextMessage() {
|
|
cursors.back()->nextMessage();
|
|
}
|
|
|
|
StringRef MultiCursor::getMessage() {
|
|
return cursors.back()->getMessage();
|
|
}
|
|
|
|
StringRef MultiCursor::getMessageWithTags() {
|
|
return cursors.back()->getMessageWithTags();
|
|
}
|
|
|
|
VectorRef<Tag> MultiCursor::getTags() const {
|
|
return cursors.back()->getTags();
|
|
}
|
|
|
|
Future<Void> MultiCursor::getMore(TaskPriority taskID) {
|
|
LogMessageVersion startVersion = cursors.back()->version();
|
|
while (cursors.size() > 1 && cursors.back()->version() >= epochEnds.back()) {
|
|
poppedVersion = std::max(poppedVersion, cursors.back()->popped());
|
|
cursors.pop_back();
|
|
epochEnds.pop_back();
|
|
}
|
|
if (cursors.back()->version() > startVersion) {
|
|
return Void();
|
|
}
|
|
return cursors.back()->getMore(taskID);
|
|
}
|
|
|
|
bool MultiCursor::isExhausted() const {
|
|
return cursors.back()->isExhausted();
|
|
}
|
|
|
|
const LogMessageVersion& MultiCursor::version() const {
|
|
return cursors.back()->version();
|
|
}
|
|
|
|
Version MultiCursor::getMinKnownCommittedVersion() const {
|
|
return cursors.back()->getMinKnownCommittedVersion();
|
|
}
|
|
|
|
Version MultiCursor::popped() const {
|
|
return std::max(poppedVersion, cursors.back()->popped());
|
|
}
|
|
|
|
TEST_CASE("/NativeCDC/ReplayPeekReplyAccounting") {
|
|
auto makeServerCursor = []() {
|
|
return makeReference<ServerPeekCursor>(
|
|
Reference<AsyncVar<OptionalInterface<TLogInterface>>>(), Tag(tagLocalityCDC, 0), 0, 100, false, false);
|
|
};
|
|
|
|
std::vector<Reference<ServerPeekCursor>> twoServers{ makeServerCursor(), makeServerCursor() };
|
|
std::vector<Reference<ServerPeekCursor>> threeServers{ makeServerCursor(), makeServerCursor(), makeServerCursor() };
|
|
auto twoReplies = makeReference<MergedPeekCursor>(twoServers, 0);
|
|
auto threeReplies = makeReference<MergedPeekCursor>(threeServers, 0);
|
|
ASSERT_EQ(twoReplies->getMaxRetainedReplyCount(), 2);
|
|
ASSERT_EQ(threeReplies->getMaxRetainedReplyCount(), 3);
|
|
|
|
std::vector<Reference<IReplayPeekCursor>> epochs{ twoReplies, threeReplies };
|
|
auto noPrefetch = makeReference<ReplayMultiCursor>(epochs, std::vector<LogMessageVersion>{ 50 }, false);
|
|
ASSERT_EQ(noPrefetch->getMaxRetainedReplyCount(), 3);
|
|
noPrefetch->setReplyByteLimit(4096);
|
|
for (const auto& cursor : twoServers) {
|
|
ASSERT_EQ(cursor->replyByteLimit, 4096);
|
|
}
|
|
for (const auto& cursor : threeServers) {
|
|
ASSERT_EQ(cursor->replyByteLimit, 4096);
|
|
}
|
|
return Void();
|
|
}
|
|
|
|
TEST_CASE("/NativeCDC/ReplayPeekCommittedEpochBoundary") {
|
|
auto makeServerCursor = [](Version committedVersion) {
|
|
auto cursor = makeReference<ServerPeekCursor>(
|
|
Reference<AsyncVar<OptionalInterface<TLogInterface>>>(), Tag(tagLocalityCDC, 0), 0, 100, false, false);
|
|
cursor->results.minKnownCommittedVersion = committedVersion;
|
|
return cursor;
|
|
};
|
|
|
|
auto current =
|
|
makeReference<MergedPeekCursor>(std::vector<Reference<ServerPeekCursor>>{ makeServerCursor(100) }, 0);
|
|
auto completed =
|
|
makeReference<MergedPeekCursor>(std::vector<Reference<ServerPeekCursor>>{ makeServerCursor(10) }, 0);
|
|
auto replay = makeReference<ReplayMultiCursor>(std::vector<Reference<IReplayPeekCursor>>{ current, completed },
|
|
std::vector<LogMessageVersion>{ LogMessageVersion(50) },
|
|
false);
|
|
|
|
ASSERT_EQ(replay->getMinKnownCommittedVersion(), 49);
|
|
replay->advanceTo(LogMessageVersion(50));
|
|
ASSERT_EQ(replay->getMinKnownCommittedVersion(), 100);
|
|
return Void();
|
|
}
|
|
|
|
BufferedCursor::BufferedCursor(std::vector<Reference<IPeekCursor>> cursors,
|
|
Version begin,
|
|
Version end,
|
|
bool withTags,
|
|
bool canDiscardPopped)
|
|
: cursors(cursors), messageIndex(0), messageVersion(begin), end(end), hasNextMessage(false), withTags(withTags),
|
|
knownUnique(false), minKnownCommittedVersion(0), poppedVersion(0), initialPoppedVersion(0),
|
|
canDiscardPopped(canDiscardPopped), randomID(deterministicRandom()->randomUniqueID()) {
|
|
ASSERT(!canDiscardPopped);
|
|
targetQueueSize = SERVER_KNOBS->DESIRED_OUTSTANDING_MESSAGES / cursors.size();
|
|
messages.reserve(SERVER_KNOBS->DESIRED_OUTSTANDING_MESSAGES);
|
|
cursorMessages.resize(cursors.size());
|
|
}
|
|
|
|
BufferedCursor::BufferedCursor(std::vector<Reference<IReplayPeekCursor>> replayCursors,
|
|
Version begin,
|
|
Version end,
|
|
bool withTags,
|
|
bool canDiscardPopped)
|
|
: discardableCursors(canDiscardPopped ? replayCursors : std::vector<Reference<IReplayPeekCursor>>()), messageIndex(0),
|
|
messageVersion(begin), end(end), hasNextMessage(false), withTags(withTags), knownUnique(false),
|
|
minKnownCommittedVersion(0), poppedVersion(0), initialPoppedVersion(0), canDiscardPopped(canDiscardPopped),
|
|
randomID(deterministicRandom()->randomUniqueID()) {
|
|
cursors.reserve(replayCursors.size());
|
|
for (const auto& cursor : replayCursors) {
|
|
cursors.push_back(cursor);
|
|
}
|
|
targetQueueSize = SERVER_KNOBS->DESIRED_OUTSTANDING_MESSAGES / cursors.size();
|
|
messages.reserve(SERVER_KNOBS->DESIRED_OUTSTANDING_MESSAGES);
|
|
cursorMessages.resize(cursors.size());
|
|
}
|
|
|
|
BufferedCursor::BufferedCursor(std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> const& logServers,
|
|
Tag tag,
|
|
Version begin,
|
|
Version end,
|
|
bool parallelGetMore)
|
|
: messageIndex(0), messageVersion(begin), end(end), hasNextMessage(false), withTags(true), knownUnique(true),
|
|
minKnownCommittedVersion(0), poppedVersion(0), initialPoppedVersion(0), canDiscardPopped(false),
|
|
randomID(deterministicRandom()->randomUniqueID()) {
|
|
targetQueueSize = SERVER_KNOBS->DESIRED_OUTSTANDING_MESSAGES / logServers.size();
|
|
messages.reserve(SERVER_KNOBS->DESIRED_OUTSTANDING_MESSAGES);
|
|
cursorMessages.resize(logServers.size());
|
|
for (int i = 0; i < logServers.size(); i++) {
|
|
auto cursor = makeReference<ServerPeekCursor>(logServers[i], tag, begin, end, false, parallelGetMore);
|
|
cursors.push_back(cursor);
|
|
}
|
|
}
|
|
|
|
void BufferedCursor::setProtocolVersion(ProtocolVersion version) {
|
|
for (auto& c : cursors) {
|
|
c->setProtocolVersion(version);
|
|
}
|
|
}
|
|
|
|
Arena& BufferedCursor::arena() {
|
|
return messages[messageIndex].arena;
|
|
}
|
|
|
|
ArenaReader* BufferedCursor::reader() {
|
|
ASSERT(false);
|
|
return cursors[0]->reader();
|
|
}
|
|
|
|
bool BufferedCursor::hasMessage() const {
|
|
return hasNextMessage;
|
|
}
|
|
|
|
void BufferedCursor::nextMessage() {
|
|
messageIndex++;
|
|
if (messageIndex == messages.size()) {
|
|
hasNextMessage = false;
|
|
}
|
|
}
|
|
|
|
StringRef BufferedCursor::getMessage() {
|
|
ASSERT(!withTags);
|
|
return messages[messageIndex].message;
|
|
}
|
|
|
|
StringRef BufferedCursor::getMessageWithTags() {
|
|
ASSERT(withTags);
|
|
return messages[messageIndex].message;
|
|
}
|
|
|
|
VectorRef<Tag> BufferedCursor::getTags() const {
|
|
ASSERT(withTags);
|
|
return messages[messageIndex].tags;
|
|
}
|
|
|
|
Future<Void> bufferedGetMoreLoader(BufferedCursor* self, Reference<IPeekCursor> cursor, int idx, TaskPriority taskID) {
|
|
while (true) {
|
|
co_await yield();
|
|
if (cursor->version().version >= self->end || self->cursorMessages[idx].size() > self->targetQueueSize) {
|
|
co_return;
|
|
}
|
|
co_await cursor->getMore(taskID);
|
|
self->poppedVersion = std::max(self->poppedVersion, cursor->popped());
|
|
self->minKnownCommittedVersion =
|
|
std::max(self->minKnownCommittedVersion, cursor->getMinKnownCommittedVersion());
|
|
if (self->canDiscardPopped) {
|
|
self->initialPoppedVersion = std::max(self->initialPoppedVersion, cursor->popped());
|
|
}
|
|
if (cursor->version().version >= self->end) {
|
|
co_return;
|
|
}
|
|
while (cursor->hasMessage()) {
|
|
self->cursorMessages[idx].push_back(
|
|
BufferedCursor::BufferedMessage(cursor->arena(),
|
|
!self->withTags ? cursor->getMessage() : cursor->getMessageWithTags(),
|
|
!self->withTags ? VectorRef<Tag>() : cursor->getTags(),
|
|
cursor->version()));
|
|
cursor->nextMessage();
|
|
}
|
|
}
|
|
}
|
|
|
|
Future<Void> bufferedGetMore(BufferedCursor* self, TaskPriority taskID) {
|
|
if (self->messageVersion.version >= self->end) {
|
|
co_await Future<Void>(Never());
|
|
throw internal_error();
|
|
}
|
|
|
|
self->messages.clear();
|
|
|
|
std::vector<Future<Void>> loaders;
|
|
loaders.reserve(self->cursors.size());
|
|
|
|
for (int i = 0; i < self->cursors.size(); i++) {
|
|
loaders.push_back(bufferedGetMoreLoader(self, self->cursors[i], i, taskID));
|
|
}
|
|
|
|
Future<Void> allLoaders = waitForAll(loaders);
|
|
Version minVersion{ 0 };
|
|
while (true) {
|
|
co_await (allLoaders || delay(SERVER_KNOBS->DESIRED_GET_MORE_DELAY, taskID));
|
|
minVersion = self->end;
|
|
for (int i = 0; i < self->cursors.size(); i++) {
|
|
auto cursor = self->cursors[i];
|
|
while (cursor->hasMessage()) {
|
|
self->cursorMessages[i].push_back(BufferedCursor::BufferedMessage(
|
|
cursor->arena(),
|
|
!self->withTags ? cursor->getMessage() : cursor->getMessageWithTags(),
|
|
!self->withTags ? VectorRef<Tag>() : cursor->getTags(),
|
|
cursor->version()));
|
|
cursor->nextMessage();
|
|
}
|
|
minVersion = std::min(minVersion, cursor->version().version);
|
|
}
|
|
if (minVersion > self->messageVersion.version) {
|
|
break;
|
|
}
|
|
if (allLoaders.isReady()) {
|
|
co_await Future<Void>(Never());
|
|
}
|
|
}
|
|
co_await yield();
|
|
|
|
for (auto& it : self->cursorMessages) {
|
|
while (!it.empty() && it.front().version.version < minVersion) {
|
|
self->messages.push_back(it.front());
|
|
it.pop_front();
|
|
}
|
|
}
|
|
if (self->knownUnique) {
|
|
std::sort(self->messages.begin(), self->messages.end());
|
|
} else {
|
|
uniquify(self->messages);
|
|
}
|
|
|
|
self->messageVersion = LogMessageVersion(minVersion);
|
|
self->messageIndex = 0;
|
|
self->hasNextMessage = !self->messages.empty();
|
|
|
|
co_await yield();
|
|
if (self->canDiscardPopped && self->poppedVersion > self->version().version) {
|
|
TraceEvent(SevWarn, "DiscardingPoppedData", self->randomID)
|
|
.detail("Version", self->version().version)
|
|
.detail("Popped", self->poppedVersion);
|
|
self->messageVersion = std::max(self->messageVersion, LogMessageVersion(self->poppedVersion));
|
|
for (const auto& cursor : self->discardableCursors) {
|
|
cursor->advanceTo(self->messageVersion);
|
|
}
|
|
self->messageIndex = self->messages.size();
|
|
if (!self->messages.empty() &&
|
|
self->messages[self->messages.size() - 1].version.version < self->poppedVersion) {
|
|
self->hasNextMessage = false;
|
|
} else {
|
|
auto iter = std::lower_bound(
|
|
self->messages.begin(), self->messages.end(), BufferedCursor::BufferedMessage(self->poppedVersion));
|
|
self->hasNextMessage = iter != self->messages.end();
|
|
if (self->hasNextMessage) {
|
|
self->messageIndex = iter - self->messages.begin();
|
|
}
|
|
}
|
|
}
|
|
if (self->hasNextMessage) {
|
|
self->canDiscardPopped = false;
|
|
}
|
|
}
|
|
|
|
Future<Void> BufferedCursor::getMore(TaskPriority taskID) {
|
|
if (hasMessage()) {
|
|
return Void();
|
|
}
|
|
|
|
if (!more.isValid() || more.isReady()) {
|
|
more = bufferedGetMore(this, taskID);
|
|
}
|
|
return more;
|
|
}
|
|
|
|
bool BufferedCursor::isExhausted() const {
|
|
ASSERT(false);
|
|
return false;
|
|
}
|
|
|
|
const LogMessageVersion& BufferedCursor::version() const {
|
|
if (hasNextMessage) {
|
|
return messages[messageIndex].version;
|
|
}
|
|
return messageVersion;
|
|
}
|
|
|
|
Version BufferedCursor::getMinKnownCommittedVersion() const {
|
|
return minKnownCommittedVersion;
|
|
}
|
|
|
|
Version BufferedCursor::popped() const {
|
|
if (initialPoppedVersion == poppedVersion) {
|
|
return 0;
|
|
}
|
|
return poppedVersion;
|
|
}
|