foundationdb/fdbserver/logsystem/LogSystemTypes.h

351 lines
14 KiB
C++

/*
* LogSystemTypes.h
*
* 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.
*/
#ifndef FDBSERVER_LOGSYSTEM_LOGSYSTEMTYPES_H
#define FDBSERVER_LOGSYSTEM_LOGSYSTEMTYPES_H
#pragma once
#include "fdbserver/logsystem/LogSystem.h"
// Leaf replay cursor backed by a single TLog interface.
class ServerPeekCursor final : public IReplayPeekCursor, public ReferenceCounted<ServerPeekCursor> {
public:
Reference<AsyncVar<OptionalInterface<TLogInterface>>> interf;
const Tag tag;
TLogPeekReply results;
ArenaReader rd;
LogMessageVersion messageVersion, end;
Version poppedVersion;
TagsAndMessage messageAndTags;
bool hasMsg;
Future<Void> more;
UID randomID;
bool returnIfBlocked;
bool onlySpilled;
bool parallelGetMore;
bool usePeekStream;
int sequence;
Deque<Future<TLogPeekReply>> futureResults;
Future<Void> interfaceChanged;
Optional<ReplyPromiseStream<TLogPeekStreamReply>> peekReplyStream;
double lastReset;
Future<Void> resetCheck;
int slowReplies;
int fastReplies;
int unknownReplies;
bool returnEmptyIfStopped;
int replyByteLimit;
ServerPeekCursor(Reference<AsyncVar<OptionalInterface<TLogInterface>>> const& interf,
Tag tag,
Version begin,
Version end,
bool returnIfBlocked,
bool parallelGetMore,
bool returnEmtpyIfStopped = false);
ServerPeekCursor(TLogPeekReply const& results,
LogMessageVersion const& messageVersion,
LogMessageVersion const& end,
TagsAndMessage const& message,
bool hasMsg,
Version poppedVersion,
Tag tag);
Reference<ServerPeekCursor> cloneServerNoMore();
Reference<IReplayPeekCursor> cloneNoMore() override;
void setProtocolVersion(ProtocolVersion version) override;
Arena& arena() override;
ArenaReader* reader() override;
bool hasMessage() const override;
void nextMessage() override;
StringRef getMessage() override;
StringRef getMessageWithTags() override;
VectorRef<Tag> getTags() const override;
void advanceTo(LogMessageVersion n) override;
Future<Void> getMore(TaskPriority taskID = TaskPriority::TLogPeekReply) override;
Future<Void> onFailed() const;
bool isActive() const;
bool isExhausted() const override;
const LogMessageVersion& version() const override;
Version popped() const override;
Version getMinKnownCommittedVersion() const override;
int64_t getMaxRetainedReplyCount() const override;
void setReplyByteLimit(int limitBytes) override;
Optional<UID> getPrimaryPeekLocation() const override;
Optional<UID> getCurrentPeekLocation() const override;
void addref() override { ReferenceCounted<ServerPeekCursor>::addref(); }
void delref() override { ReferenceCounted<ServerPeekCursor>::delref(); }
Version getMaxKnownVersion() const override { return results.maxKnownVersion; }
};
// Replay cursor that reads one logical stream from replicated TLog servers in a log set.
class MergedPeekCursor final : public IReplayPeekCursor, public ReferenceCounted<MergedPeekCursor> {
public:
Reference<LogSet> logSet;
std::vector<Reference<ServerPeekCursor>> serverCursors;
std::vector<LocalityEntry> locations;
std::vector<std::pair<LogMessageVersion, int>> sortedVersions;
Tag tag;
int bestServer, currentCursor, readQuorum;
Optional<LogMessageVersion> nextVersion;
LogMessageVersion messageVersion;
bool hasNextMessage;
UID randomID;
int tLogReplicationFactor;
Future<Void> more;
MergedPeekCursor(std::vector<Reference<ServerPeekCursor>> const& serverCursors, Version begin);
MergedPeekCursor(std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> const& logServers,
int bestServer,
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 = Optional<std::vector<uint16_t>>());
MergedPeekCursor(std::vector<Reference<ServerPeekCursor>> const& serverCursors,
LogMessageVersion const& messageVersion,
int bestServer,
int readQuorum,
Optional<LogMessageVersion> nextVersion,
Reference<LogSet> logSet,
int tLogReplicationFactor);
Reference<IReplayPeekCursor> cloneNoMore() override;
void setProtocolVersion(ProtocolVersion version) override;
Arena& arena() override;
ArenaReader* reader() override;
void calcHasMessage();
void updateMessage(bool usePolicy);
bool hasMessage() const override;
void nextMessage() override;
StringRef getMessage() override;
StringRef getMessageWithTags() override;
VectorRef<Tag> getTags() const override;
void advanceTo(LogMessageVersion n) override;
Future<Void> getMore(TaskPriority taskID = TaskPriority::TLogPeekReply) override;
bool isExhausted() const override;
const LogMessageVersion& version() const override;
Version popped() const override;
Version getMinKnownCommittedVersion() const override;
Version getMaxKnownVersion() const override;
int64_t getMaxRetainedReplyCount() const override;
void setReplyByteLimit(int limitBytes) override;
Optional<UID> getPrimaryPeekLocation() const override;
Optional<UID> getCurrentPeekLocation() const override;
void addref() override { ReferenceCounted<MergedPeekCursor>::addref(); }
void delref() override { ReferenceCounted<MergedPeekCursor>::delref(); }
};
// Replay cursor that reads one logical stream across candidate log sets.
class SetPeekCursor final : public IReplayPeekCursor, public ReferenceCounted<SetPeekCursor> {
public:
std::vector<Reference<LogSet>> logSets;
std::vector<std::vector<Reference<ServerPeekCursor>>> serverCursors;
Tag tag;
int bestSet, bestServer, currentSet, currentCursor;
std::vector<LocalityEntry> locations;
std::vector<std::pair<LogMessageVersion, int>> sortedVersions;
Optional<LogMessageVersion> nextVersion;
LogMessageVersion messageVersion;
bool hasNextMessage;
bool useBestSet;
UID randomID;
Future<Void> more;
Optional<Version> end;
SetPeekCursor(std::vector<Reference<LogSet>> const& logSets,
int bestSet,
int bestServer,
Tag tag,
Version begin,
Version end,
bool parallelGetMore,
const Optional<std::vector<uint16_t>>& knownLockedTLogIds = Optional<std::vector<uint16_t>>());
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);
Reference<IReplayPeekCursor> cloneNoMore() override;
void setProtocolVersion(ProtocolVersion version) override;
Arena& arena() override;
ArenaReader* reader() override;
void calcHasMessage();
void updateMessage(int logIdx, bool usePolicy);
bool hasMessage() const override;
void nextMessage() override;
StringRef getMessage() override;
StringRef getMessageWithTags() override;
VectorRef<Tag> getTags() const override;
void advanceTo(LogMessageVersion n) override;
Future<Void> getMore(TaskPriority taskID = TaskPriority::TLogPeekReply) override;
bool isExhausted() const override;
const LogMessageVersion& version() const override;
Version popped() const override;
Version getMinKnownCommittedVersion() const override;
Version getMaxKnownVersion() const override;
int64_t getMaxRetainedReplyCount() const override;
void setReplyByteLimit(int limitBytes) override;
Optional<UID> getPrimaryPeekLocation() const override;
Optional<UID> getCurrentPeekLocation() const override;
void addref() override { ReferenceCounted<SetPeekCursor>::addref(); }
void delref() override { ReferenceCounted<SetPeekCursor>::delref(); }
};
// Replay cursor that stitches together replay-capable cursors from successive history ranges.
class ReplayMultiCursor final : public IReplayPeekCursor, public ReferenceCounted<ReplayMultiCursor> {
public:
std::vector<Reference<IReplayPeekCursor>> cursors;
std::vector<LogMessageVersion> epochEnds;
Version poppedVersion;
bool prefetch;
ReplayMultiCursor(std::vector<Reference<IReplayPeekCursor>> cursors,
std::vector<LogMessageVersion> epochEnds,
bool prefetch = true);
Reference<IReplayPeekCursor> cloneNoMore() override;
void setProtocolVersion(ProtocolVersion version) override;
Arena& arena() override;
ArenaReader* reader() override;
bool hasMessage() const override;
void nextMessage() override;
StringRef getMessage() override;
StringRef getMessageWithTags() override;
VectorRef<Tag> getTags() const override;
void advanceTo(LogMessageVersion n) override;
Future<Void> getMore(TaskPriority taskID = TaskPriority::TLogPeekReply) override;
bool isExhausted() const override;
const LogMessageVersion& version() const override;
Version popped() const override;
Version getMinKnownCommittedVersion() const override;
Version getMaxKnownVersion() const override;
int64_t getMaxRetainedReplyCount() const override;
void setReplyByteLimit(int limitBytes) override;
Optional<UID> getPrimaryPeekLocation() const override;
Optional<UID> getCurrentPeekLocation() const override;
void addref() override { ReferenceCounted<ReplayMultiCursor>::addref(); }
void delref() override { ReferenceCounted<ReplayMultiCursor>::delref(); }
};
// Plain cursor that stitches together sequential ranges without replay-only capabilities.
class MultiCursor final : public IPeekCursor, public ReferenceCounted<MultiCursor> {
public:
std::vector<Reference<IPeekCursor>> cursors;
std::vector<LogMessageVersion> epochEnds;
Version poppedVersion;
MultiCursor(std::vector<Reference<IPeekCursor>> cursors, std::vector<LogMessageVersion> epochEnds);
void setProtocolVersion(ProtocolVersion version) override;
Arena& arena() override;
ArenaReader* reader() override;
bool hasMessage() const override;
void nextMessage() override;
StringRef getMessage() override;
StringRef getMessageWithTags() override;
VectorRef<Tag> getTags() const override;
Future<Void> getMore(TaskPriority taskID = TaskPriority::TLogPeekReply) override;
bool isExhausted() const override;
const LogMessageVersion& version() const override;
Version popped() const override;
Version getMinKnownCommittedVersion() const override;
void addref() override { ReferenceCounted<MultiCursor>::addref(); }
void delref() override { ReferenceCounted<MultiCursor>::delref(); }
};
// Plain cursor that buffers and orders messages pulled from one or more input cursors.
class BufferedCursor final : public IPeekCursor, public ReferenceCounted<BufferedCursor> {
public:
struct BufferedMessage {
Arena arena;
StringRef message;
VectorRef<Tag> tags;
LogMessageVersion version;
BufferedMessage() = default;
explicit BufferedMessage(Version version) : version(version) {}
BufferedMessage(Arena arena, StringRef message, const VectorRef<Tag>& tags, const LogMessageVersion& version)
: arena(arena), message(message), tags(tags), version(version) {}
bool operator<(BufferedMessage const& r) const { return version < r.version; }
bool operator==(BufferedMessage const& r) const { return version == r.version; }
};
std::vector<Reference<IPeekCursor>> cursors;
std::vector<Reference<IReplayPeekCursor>> discardableCursors;
std::vector<Deque<BufferedMessage>> cursorMessages;
std::vector<BufferedMessage> messages;
int messageIndex;
LogMessageVersion messageVersion;
Version end;
bool hasNextMessage;
bool withTags;
bool knownUnique;
Version minKnownCommittedVersion;
Version poppedVersion;
Version initialPoppedVersion;
bool canDiscardPopped;
Future<Void> more;
int targetQueueSize;
UID randomID;
BufferedCursor(std::vector<Reference<IPeekCursor>> cursors,
Version begin,
Version end,
bool withTags,
bool canDiscardPopped);
BufferedCursor(std::vector<Reference<IReplayPeekCursor>> cursors,
Version begin,
Version end,
bool withTags,
bool canDiscardPopped);
BufferedCursor(std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> const& logServers,
Tag tag,
Version begin,
Version end,
bool parallelGetMore);
void setProtocolVersion(ProtocolVersion version) override;
Arena& arena() override;
ArenaReader* reader() override;
bool hasMessage() const override;
void nextMessage() override;
StringRef getMessage() override;
StringRef getMessageWithTags() override;
VectorRef<Tag> getTags() const override;
Future<Void> getMore(TaskPriority taskID = TaskPriority::TLogPeekReply) override;
bool isExhausted() const override;
const LogMessageVersion& version() const override;
Version popped() const override;
Version getMinKnownCommittedVersion() const override;
void addref() override { ReferenceCounted<BufferedCursor>::addref(); }
void delref() override { ReferenceCounted<BufferedCursor>::delref(); }
};
#endif