351 lines
14 KiB
C++
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
|