diff --git a/dCommon/CMakeLists.txt b/dCommon/CMakeLists.txt index 3fb9d4d32..62cca5638 100644 --- a/dCommon/CMakeLists.txt +++ b/dCommon/CMakeLists.txt @@ -4,6 +4,7 @@ set(DCOMMON_SOURCES "BinaryIO.cpp" "dConfig.cpp" "Diagnostics.cpp" + "TrafficStats.cpp" "Locale.cpp" "Logger.cpp" "Game.cpp" diff --git a/dCommon/TrafficStats.cpp b/dCommon/TrafficStats.cpp new file mode 100644 index 000000000..fd79da820 --- /dev/null +++ b/dCommon/TrafficStats.cpp @@ -0,0 +1,276 @@ +#include "TrafficStats.h" + +#include +#include +#include + +#include "MessageIdentifiers.h" +#include "ServiceType.h" +#include "MessageType/Client.h" +#include "MessageType/World.h" + +namespace TrafficStats { + namespace { + constexpr double FIRST_BOUND = 100.0; // microseconds + + const std::array& Bounds() { + static const auto bounds = [] { + std::array out{}; + for (size_t i = 0; i + 1 < Histogram::BUCKETS; i++) out[i] = static_cast(std::llround(FIRST_BOUND * std::exp2(static_cast(i) / 3.0))); + out[Histogram::BUCKETS - 1] = UINT64_MAX; + return out; + }(); + return bounds; + } + + template + void AddArrays(std::array& to, const std::array& from) { + for (size_t i = 0; i < N; i++) to[i] += from[i]; + } + } + + uint64_t Histogram::UpperBound(size_t bucket) { + return Bounds()[std::min(bucket, BUCKETS - 1)]; + } + + size_t Histogram::BucketFor(uint64_t microseconds) { + const auto& bounds = Bounds(); + return static_cast(std::lower_bound(bounds.begin(), bounds.end(), microseconds) - bounds.begin()); + } + + void Histogram::Add(uint64_t microseconds, uint32_t count) { + m_Counts[BucketFor(microseconds)] += count; + m_Count += count; + m_Sum += microseconds * count; + } + + void Histogram::AddBucket(size_t bucket, uint32_t count) { + if (bucket >= BUCKETS) return; + m_Counts[bucket] += count; + m_Count += count; + } + + void Histogram::Merge(const Histogram& other) { + for (size_t i = 0; i < BUCKETS; i++) m_Counts[i] += other.m_Counts[i]; + m_Count += other.m_Count; + m_Sum += other.m_Sum; + } + + uint64_t Histogram::Percentile(double fraction) const { + if (m_Count == 0) return 0; + fraction = std::clamp(fraction, 0.0, 1.0); + // The rank of the value wanted, 1-based: the smallest value is rank 1 + const double rank = std::max(1.0, std::ceil(fraction * static_cast(m_Count))); + uint64_t seen = 0; + for (size_t i = 0; i < BUCKETS; i++) { + if (m_Counts[i] == 0) continue; + if (static_cast(seen + m_Counts[i]) >= rank) { + const double lower = i == 0 ? 0.0 : static_cast(UpperBound(i - 1)); + // The overflow bucket has no upper bound: report its lower one + if (i == BUCKETS - 1) return static_cast(lower); + const double upper = static_cast(UpperBound(i)); + const double within = (rank - static_cast(seen)) / static_cast(m_Counts[i]); + return static_cast(std::llround(lower + (upper - lower) * within)); + } + seen += m_Counts[i]; + } + return 0; + } + + std::vector> Histogram::Sparse() const { + std::vector> out; + for (size_t i = 0; i < BUCKETS; i++) { + if (m_Counts[i]) out.emplace_back(static_cast(i), m_Counts[i]); + } + return out; + } + + Histogram Histogram::FromSparse(const std::vector>& sparse, uint64_t sum) { + Histogram out; + for (const auto& [bucket, count] : sparse) out.AddBucket(bucket, count); + out.m_Sum = sum; + return out; + } + + size_t StatusClass(uint16_t status) { + if (status >= 100 && status < 500) return status / 100 - 1; + return 4; + } + + uint64_t MessageKey::Packed() const { + return (static_cast(outbound) << 63) | (static_cast(service & 0x7FFF) << 48) | + (static_cast(packet) << 16) | gameMessage; + } + + MessageKey MessageKey::Unpack(uint64_t packed) { + MessageKey key; + key.outbound = (packed >> 63) != 0; + key.service = static_cast((packed >> 48) & 0x7FFF); + if (key.service == (RAKNET & 0x7FFF)) key.service = RAKNET; + key.packet = static_cast(packed >> 16); + key.gameMessage = static_cast(packed); + return key; + } + + MessageKey KeyOf(const uint8_t* data, size_t length, bool outbound) { + MessageKey key; + key.outbound = outbound; + if (!data || length == 0) { + key.service = MessageKey::RAKNET; + return key; + } + // LU packets: ID_USER_PACKET_ENUM, uint16 service, uint32 packet ID, one padding byte + if (data[0] != ID_USER_PACKET_ENUM || length < 8) { + key.service = MessageKey::RAKNET; + key.packet = data[0]; + return key; + } + key.service = static_cast(data[1] | (data[2] << 8)); + key.packet = static_cast(data[3]) | (static_cast(data[4]) << 8) | (static_cast(data[5]) << 16) | (static_cast(data[6]) << 24); + // Game messages: the header, the target object (8 bytes), then the uint16 game message ID + const bool gameMessage = (key.service == static_cast(ServiceType::WORLD) && key.packet == static_cast(MessageType::World::GAME_MSG)) || + (key.service == static_cast(ServiceType::CLIENT) && key.packet == static_cast(MessageType::Client::GAME_MSG)); + if (gameMessage && length >= 18) key.gameMessage = static_cast(data[16] | (data[17] << 8)); + return key; + } + + void Second::Merge(const Second& other) { + packetsIn += other.packetsIn; + packetsOut += other.packetsOut; + bytesIn += other.bytesIn; + bytesOut += other.bytesOut; + httpRequests += other.httpRequests; + AddArrays(httpStatus, other.httpStatus); + httpBytesOut += other.httpBytesOut; + httpLatency.Merge(other.httpLatency); + } + + void RouteStats::Merge(const RouteStats& other) { + count += other.count; + AddArrays(status, other.status); + bytesOut += other.bytesOut; + latency.Merge(other.latency); + } + + void Recorder::Packet(int64_t now, const MessageKey& key, uint64_t bytes, uint32_t fanout) { + if (fanout == 0) return; + std::lock_guard lock(m_Mutex); + auto& second = SecondAt(now); + if (key.outbound) { + second.packetsOut += fanout; + second.bytesOut += bytes * fanout; + } else { + second.packetsIn += fanout; + second.bytesIn += bytes * fanout; + } + auto& message = m_Messages[key.Packed()]; + message.key = key; + message.count += fanout; + message.bytes += bytes * fanout; + } + + void Recorder::Http(int64_t now, const std::string& route, uint16_t status, uint64_t microseconds, uint64_t bytesOut) { + std::lock_guard lock(m_Mutex); + auto& second = SecondAt(now); + second.httpRequests++; + second.httpStatus[StatusClass(status)]++; + second.httpBytesOut += bytesOut; + second.httpLatency.Add(microseconds); + + auto it = m_Routes.find(route); + if (it == m_Routes.end()) { + const bool full = m_Routes.size() >= MAX_ROUTES; + it = m_Routes.try_emplace(full ? std::string("other") : route).first; + it->second.route = it->first; + } + it->second.count++; + it->second.status[StatusClass(status)]++; + it->second.bytesOut += bytesOut; + it->second.latency.Add(microseconds); + } + + Second& Recorder::SecondAt(int64_t now) { + // Nearly every packet falls in the same second as the one before + if (m_Current && m_CurrentTime == now) return *m_Current; + auto& second = m_Seconds[now]; + second.time = now; + m_Current = &second; + m_CurrentTime = now; + return second; + } + + void Recorder::SetGauge(const std::string& name, std::function source) { + std::lock_guard lock(m_Mutex); + for (auto& gauge : m_Gauges) { + if (gauge.first == name) { + gauge.second = std::move(source); + return; + } + } + m_Gauges.emplace_back(name, std::move(source)); + } + + bool Recorder::Due(int64_t now, int64_t interval) { + std::lock_guard lock(m_Mutex); + if (m_LastTake == 0) { + m_LastTake = now; + return false; + } + return now - m_LastTake >= interval; + } + + Report Recorder::Take(int64_t now) { + Report report; + std::vector>> gauges; + { + std::lock_guard lock(m_Mutex); + m_LastTake = now; + // Fill every second from the last report to the one before now; after a long silence only the last MAX_GAP + int64_t from = m_LastReported ? m_LastReported + 1 : (m_Seconds.empty() ? now : std::min(m_Seconds.begin()->first, now - 1)); + from = std::max(from, now - MAX_GAP); + for (int64_t t = from; t < now; t++) { + const auto it = m_Seconds.find(t); + if (it != m_Seconds.end()) report.seconds.push_back(std::move(it->second)); + else report.seconds.push_back(Second{ .time = t }); + } + m_Seconds.erase(m_Seconds.begin(), m_Seconds.lower_bound(now)); + m_Current = nullptr; + if (now - 1 > m_LastReported) m_LastReported = now - 1; + + std::vector messages; + messages.reserve(m_Messages.size()); + for (auto& [_, count] : m_Messages) messages.push_back(count); + m_Messages.clear(); + report.messages = Top(messages, TOP_MESSAGES); + + for (auto& [_, route] : m_Routes) report.routes.push_back(std::move(route)); + m_Routes.clear(); + gauges = m_Gauges; + } + for (const auto& [name, source] : gauges) report.gauges.emplace_back(name, source ? source() : 0.0); + return report; + } + + std::vector Top(const std::vector& counts, size_t limit) { + std::vector out; + for (const bool outbound : { false, true }) { + std::vector direction; + for (const auto& count : counts) if (count.key.outbound == outbound) direction.push_back(count); + std::sort(direction.begin(), direction.end(), [](const MessageCount& a, const MessageCount& b) { + return a.count != b.count ? a.count > b.count : a.key.Packed() < b.key.Packed(); + }); + if (direction.size() > limit) direction.resize(limit); + out.insert(out.end(), direction.begin(), direction.end()); + } + return out; + } + + Recorder& Local() { + static Recorder recorder; + return recorder; + } + + int64_t Now() { + return static_cast(std::time(nullptr)); + } +} diff --git a/dCommon/TrafficStats.h b/dCommon/TrafficStats.h new file mode 100644 index 000000000..36e9dc560 --- /dev/null +++ b/dCommon/TrafficStats.h @@ -0,0 +1,173 @@ +#pragma once + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +/** + * Traffic diagnostics every server keeps about itself: packets and bytes in and out per second, the packet and game + * message types they were, and (servers with a web server) HTTP requests with their latency. Counting is a few + * additions under an uncontended lock; every few seconds dServer takes a Report and ships it to the dashboard through + * master (see docs/Dashboard.md, "Traffic diagnostics"). + * + * Seconds are Unix seconds, so reports from different servers line up. + */ +namespace TrafficStats { + /** + * Latency histogram with fixed, geometric buckets (three per doubling, from 100 microseconds), so histograms from + * different seconds and servers add up and percentiles of a minute or an hour come out right. + */ + class Histogram { + public: + static constexpr size_t BUCKETS = 58; // the last one holds everything over ~41 s + + // Largest value (microseconds) bucket i holds; the last bucket has no limit (UINT64_MAX) + static uint64_t UpperBound(size_t bucket); + static size_t BucketFor(uint64_t microseconds); + + void Add(uint64_t microseconds, uint32_t count = 1); + void Merge(const Histogram& other); + void AddBucket(size_t bucket, uint32_t count); + + uint64_t Count() const { return m_Count; } + uint64_t Sum() const { return m_Sum; } // microseconds, all values together + void SetSum(uint64_t sum) { m_Sum = sum; } + bool Empty() const { return m_Count == 0; } + uint32_t At(size_t bucket) const { return m_Counts[bucket]; } + + // The value below which `fraction` (0..1) of the values are, interpolated within its bucket; 0 when empty + uint64_t Percentile(double fraction) const; + + // Non-empty buckets as (bucket, count), for the wire and for keeping many histograms small + std::vector> Sparse() const; + static Histogram FromSparse(const std::vector>& sparse, uint64_t sum); + + private: + std::array m_Counts{}; + uint64_t m_Count{}; + uint64_t m_Sum{}; + }; + + // HTTP status classes 1xx..5xx as index 0..4 (anything else counts as 5xx) + size_t StatusClass(uint16_t status); + + // What a packet was: its LU service and packet ID, and for game messages the game message ID + struct MessageKey { + bool outbound{}; + uint16_t service{}; // ServiceType; RAKNET for RakNet's own messages (no LU header) + uint32_t packet{}; // LU packet ID, or the RakNet message ID + uint16_t gameMessage{}; // game messages only + + static constexpr uint16_t RAKNET = 0xFFFF; + + uint64_t Packed() const; + static MessageKey Unpack(uint64_t packed); + bool operator==(const MessageKey&) const = default; + }; + + // Reads the key from a packet's first bytes (never reads past `length`) + MessageKey KeyOf(const uint8_t* data, size_t length, bool outbound); + + struct MessageCount { + MessageKey key; + uint64_t count{}; + uint64_t bytes{}; + }; + + // One second of traffic + struct Second { + int64_t time{}; + uint64_t packetsIn{}; + uint64_t packetsOut{}; + uint64_t bytesIn{}; + uint64_t bytesOut{}; + uint64_t httpRequests{}; + std::array httpStatus{}; // by StatusClass + uint64_t httpBytesOut{}; + Histogram httpLatency; + + void Merge(const Second& other); // adds the counts (keeps this one's time) + bool Idle() const { return packetsIn == 0 && packetsOut == 0 && httpRequests == 0; } + }; + + struct RouteStats { + std::string route; // "GET /api/players/:id" + uint64_t count{}; + std::array status{}; + uint64_t bytesOut{}; + Histogram latency; + + void Merge(const RouteStats& other); + }; + + // RakNet's view of the connections (all datagrams, acknowledgements and resends included), over the report + struct Link { + uint32_t connections{}; + uint64_t datagramsSent{}; + uint64_t datagramsReceived{}; + uint64_t bytesSent{}; + uint64_t bytesReceived{}; + uint64_t resends{}; + uint32_t resendQueue{}; // messages waiting to be resent now + uint32_t averagePingMs{}; // over the connections + }; + + struct Report { + std::vector seconds; // oldest first, one per second without gaps + std::vector messages; // the busiest types per direction + std::vector routes; + Link link; + std::vector> gauges; // e.g. workers_busy + }; + + class Recorder { + public: + static constexpr size_t TOP_MESSAGES = 24; // per direction, per report + static constexpr size_t MAX_ROUTES = 64; // distinct routes per report; more count as "other" + static constexpr int64_t MAX_GAP = 120; // seconds of silence a report fills in at most + + // `fanout`: how many connections a broadcast went to + void Packet(int64_t now, const MessageKey& key, uint64_t bytes, uint32_t fanout = 1); + void Http(int64_t now, const std::string& route, uint16_t status, uint64_t microseconds, uint64_t bytesOut); + + // Evaluated when a report is taken (on the thread that takes it) + void SetGauge(const std::string& name, std::function source); + + /** + * The seconds before `now` not reported yet (silent ones as zeros, at most MAX_GAP of them), the busiest + * message types and the routes since the last report. The current second stays for the next one. + */ + Report Take(int64_t now); + + // Whether `interval` seconds passed since the last Take (the first call only starts the clock) + bool Due(int64_t now, int64_t interval); + + private: + Second& SecondAt(int64_t now); // under m_Mutex + + std::mutex m_Mutex; + std::map m_Seconds; + Second* m_Current{}; // m_Seconds[m_CurrentTime] (map entries stay put) + int64_t m_CurrentTime{}; + int64_t m_LastReported{}; // last second a report covered + std::unordered_map m_Messages; + std::map m_Routes; + std::vector>> m_Gauges; + int64_t m_LastTake{}; + }; + + // The busiest `limit` message types of each direction, busiest first + std::vector Top(const std::vector& counts, size_t limit); + + // This process's recorder (dServer counts packets into it, the web server requests) + Recorder& Local(); + + int64_t Now(); // Unix seconds +} diff --git a/dCommon/dEnums/MessageType/Master.h b/dCommon/dEnums/MessageType/Master.h index 188f7b7e3..c0c531551 100644 --- a/dCommon/dEnums/MessageType/Master.h +++ b/dCommon/dEnums/MessageType/Master.h @@ -63,5 +63,8 @@ namespace MessageType { MIGRATE_STATUS, // Source world -> master -> target world: what a moved player had that isn't in their saved character MIGRATE_PLAYER_STATE, + + // Any server -> master -> dashboard: traffic counters of the last few seconds (see ServerTraffic.h) + SERVER_TRAFFIC, }; } diff --git a/dMasterServer/MasterServer.cpp b/dMasterServer/MasterServer.cpp index a3e99a13c..df6b14fcb 100644 --- a/dMasterServer/MasterServer.cpp +++ b/dMasterServer/MasterServer.cpp @@ -51,6 +51,7 @@ #include "master/DataChanged.h" #include "master/MessageCapture.h" #include "master/InstanceMigration.h" +#include "master/ServerTraffic.h" #ifdef DARKFLAME_PLATFORM_UNIX @@ -396,6 +397,10 @@ int main(int argc, char** argv) { assert(res == 0); Game::server = new dServer(ourIP, ourPort, 0, maxClients, true, false, Game::logger, "", 0, ServiceType::MASTER, Game::config, &Game::lastSignal, hash); + // Master has no master to send its traffic report to: it goes straight to the dashboard + Game::server->SetTrafficSink([](ServerTraffic& report) { + if (dashboardServerMasterPeerSysAddr != UNASSIGNED_SYSTEM_ADDRESS) MasterPackets::SendTo(dashboardServerMasterPeerSysAddr, report); + }); std::string master_server_ip = "localhost"; const auto masterServerIPString = Game::config->GetValue("master_ip"); @@ -855,6 +860,12 @@ namespace { MasterPackets::SendTo(dashboardServerMasterPeerSysAddr, msg); } + // Every server's traffic report goes on to the dashboard (the dashboard keeps its own) + void OnServerTraffic(const ServerTraffic& report, const SystemAddress& sysAddr) { + if (dashboardServerMasterPeerSysAddr == UNASSIGNED_SYSTEM_ADDRESS || sysAddr == dashboardServerMasterPeerSysAddr) return; + MasterPackets::SendTo(dashboardServerMasterPeerSysAddr, report); + } + void OnAnnounce(const Announcement& announcement, const SystemAddress& sysAddr) { if (sysAddr != dashboardServerMasterPeerSysAddr) { LOG("Ignoring announcement from a server that is not the dashboard"); @@ -975,6 +986,7 @@ namespace { handlers.On(Master::MESSAGE_CAPTURE_CONTROL, OnMessageCaptureControl); handlers.On(Master::MESSAGE_CAPTURE_DATA, ForwardWorldToDashboard); handlers.On(Master::REQUEST_SERVER_LIST, OnRequestServerList); + handlers.On(Master::SERVER_TRAFFIC, OnServerTraffic); return handlers; }(); return handlers; diff --git a/dNet/dServer.cpp b/dNet/dServer.cpp index 897885a67..5f7081b61 100644 --- a/dNet/dServer.cpp +++ b/dNet/dServer.cpp @@ -17,6 +17,9 @@ #include "ZoneInstanceManager.h" #include "StringifiedEnum.h" #include "GeneralUtils.h" +#include "TrafficStats.h" +#include "RakNetStatistics.h" +#include "master/ServerTraffic.h" //! Replica Constructor class class ReplicaConstructor : public ReceiveConstructionInterface { @@ -134,6 +137,7 @@ Packet* dServer::ReceiveFromMaster() { if (!mMasterConnectionActive) ConnectToMaster(); Packet* packet = mMasterPeer->Receive(); + CountTraffic(packet); if (packet) { if (packet->length < 1) { mMasterPeer->DeallocatePacket(packet); return nullptr; } @@ -197,7 +201,9 @@ Packet* dServer::ReceiveFromMaster() { } Packet* dServer::Receive() { - return mPeer->Receive(); + Packet* packet = mPeer->Receive(); + CountTraffic(packet); + return packet; } void dServer::DeallocatePacket(Packet* packet) { @@ -210,11 +216,13 @@ void dServer::DeallocateMasterPacket(Packet* packet) { void dServer::Send(RakNet::BitStream& bitStream, const SystemAddress& sysAddr, bool broadcast) { if (mSendObserver) mSendObserver(bitStream, sysAddr, broadcast); + CountTraffic(bitStream, broadcast, sysAddr); mPeer->Send(&bitStream, SYSTEM_PRIORITY, RELIABLE_ORDERED, 0, sysAddr, broadcast); } void dServer::SendToMaster(RakNet::BitStream& bitStream) { if (!mMasterConnectionActive) ConnectToMaster(); + CountTraffic(bitStream, false, mMasterSystemAddress); mMasterPeer->Send(&bitStream, SYSTEM_PRIORITY, RELIABLE_ORDERED, 0, mMasterSystemAddress, false); } @@ -223,6 +231,7 @@ void dServer::Disconnect(const SystemAddress& sysAddr, eServerDisconnectIdentifi notify.disconnectID = disconNotifyID; RakNet::BitStream bitStream; notify.WritePacket(bitStream); + CountTraffic(bitStream, false, sysAddr); mPeer->Send(&bitStream, SYSTEM_PRIORITY, RELIABLE_ORDERED, 0, sysAddr, false); mPeer->CloseConnection(sysAddr, true); @@ -316,3 +325,65 @@ int dServer::GetPing(const SystemAddress& sysAddr) const { int dServer::GetLatestPing(const SystemAddress& sysAddr) const { return mPeer->GetLastPing(sysAddr); } + +void dServer::CountTraffic(const Packet* packet) { + const auto now = TrafficStats::Now(); + if (packet) TrafficStats::Local().Packet(now, TrafficStats::KeyOf(packet->data, packet->length, false), packet->length); + if (TrafficStats::Local().Due(now, ServerTraffic::REPORT_SECONDS)) ReportTraffic(); +} + +void dServer::CountTraffic(const RakNet::BitStream& bitStream, bool broadcast, const SystemAddress& sysAddr) { + // A broadcast goes to every connection but the one given + uint32_t fanout = 1; + if (broadcast) { + const uint32_t connections = mPeer ? mPeer->NumberOfConnections() : 0; + fanout = sysAddr == UNASSIGNED_SYSTEM_ADDRESS ? connections : (connections > 0 ? connections - 1 : 0); + } + const auto bytes = bitStream.GetNumberOfBytesUsed(); + TrafficStats::Local().Packet(TrafficStats::Now(), TrafficStats::KeyOf(bitStream.GetData(), bytes, true), bytes, fanout); +} + +void dServer::AddLinkStats(RakPeerInterface* peer, uint64_t peerIndex, ServerTraffic& report, uint64_t& pingSum, std::map& seen) { + if (!peer) return; + std::vector addresses(std::max(peer->GetMaximumNumberOfPeers(), 1)); + unsigned short count = static_cast(addresses.size()); + if (!peer->GetConnectionList(addresses.data(), &count)) return; + auto& link = report.report.link; + for (unsigned short i = 0; i < count; i++) { + auto* stats = peer->GetStatistics(addresses[i]); + if (!stats) continue; + LinkCounters now{ stats->packetsSent, stats->packetsReceived, stats->totalBitsSent, stats->bitsReceived, stats->messageResends }; + const uint64_t key = (peerIndex << 48) | (static_cast(addresses[i].binaryAddress) << 16) | addresses[i].port; + const auto it = mLinkCounters.find(key); + const LinkCounters before = it != mLinkCounters.end() ? it->second : LinkCounters{}; + // A reused address is a new connection whose totals started again + const auto delta = [](uint64_t current, uint64_t previous) { return current >= previous ? current - previous : current; }; + link.datagramsSent += delta(now.datagramsSent, before.datagramsSent); + link.datagramsReceived += delta(now.datagramsReceived, before.datagramsReceived); + link.bytesSent += delta(now.bitsSent, before.bitsSent) / 8; + link.bytesReceived += delta(now.bitsReceived, before.bitsReceived) / 8; + link.resends += delta(now.resends, before.resends); + link.resendQueue += stats->messagesOnResendQueue; + link.connections++; + pingSum += std::max(0, peer->GetAveragePing(addresses[i])); + seen[key] = now; + } +} + +void dServer::ReportTraffic() { + ServerTraffic report; + report.serverType = mServerType; + report.zoneId = mZoneID; + report.instanceId = static_cast(mInstanceID); + report.report = TrafficStats::Local().Take(TrafficStats::Now()); + + uint64_t pingSum = 0; + std::map seen; + AddLinkStats(mPeer, 0, report, pingSum, seen); + AddLinkStats(mMasterPeer, 1, report, pingSum, seen); + mLinkCounters = std::move(seen); + if (report.report.link.connections) report.report.link.averagePingMs = static_cast(pingSum / report.report.link.connections); + + if (mTrafficSink) mTrafficSink(report); + else if (mMasterPeer && mMasterConnectionActive) MasterPackets::SendToMaster(report, this); +} diff --git a/dNet/dServer.h b/dNet/dServer.h index d280ae8fa..092dddab4 100644 --- a/dNet/dServer.h +++ b/dNet/dServer.h @@ -3,12 +3,14 @@ #include #include #include +#include #include "RakPeerInterface.h" #include "ReplicaManager.h" #include "NetworkIDManager.h" class Logger; class dConfig; +struct ServerTraffic; enum class eServerDisconnectIdentifiers : uint32_t; enum class ServiceType : uint16_t; @@ -51,6 +53,11 @@ public: void Disconnect(const SystemAddress& sysAddr, eServerDisconnectIdentifiers disconNotifyID); + // Where this server's traffic report goes every few seconds (see ServerTraffic.h). By default it is sent to + // master; master sends its own to the dashboard, the dashboard keeps its own. + using TrafficSink = std::function; + void SetTrafficSink(TrafficSink sink) { mTrafficSink = std::move(sink); } + bool IsConnected(const SystemAddress& sysAddr); const std::string& GetIP() const { return mIP; } const int GetPort() const { return mPort; } @@ -80,7 +87,16 @@ public: } private: + struct LinkCounters { + uint64_t datagramsSent{}, datagramsReceived{}, bitsSent{}, bitsReceived{}, resends{}; + }; bool Startup(); + // Traffic diagnostics (TrafficStats): count one packet, and send the report when it is due + void CountTraffic(const Packet* packet); + void CountTraffic(const RakNet::BitStream& bitStream, bool broadcast, const SystemAddress& sysAddr); + void ReportTraffic(); + // Adds the peer's connections to the report's link statistics (changes since the last report) + void AddLinkStats(RakPeerInterface* peer, uint64_t peerIndex, ServerTraffic& report, uint64_t& pingSum, std::map& seen); void Shutdown(); void SetupForMasterConnection(); bool ConnectToMaster(); @@ -118,4 +134,8 @@ protected: std::chrono::steady_clock::time_point mStartTime = std::chrono::steady_clock::now(); std::string mMasterPassword; SendObserver mSendObserver; + + TrafficSink mTrafficSink; + // RakNet's per-connection statistics are totals since the connection opened; the last ones seen, for deltas + std::map mLinkCounters; }; diff --git a/dNet/master/ServerTraffic.h b/dNet/master/ServerTraffic.h new file mode 100644 index 000000000..227db21dc --- /dev/null +++ b/dNet/master/ServerTraffic.h @@ -0,0 +1,172 @@ +#ifndef __SERVERTRAFFIC__H__ +#define __SERVERTRAFFIC__H__ + +#include +#include +#include +#include + +#include "BitStream.h" +#include "BitStreamUtils.h" +#include "MessageType/Master.h" +#include "ServiceType.h" +#include "TrafficStats.h" + +/** + * SERVER_TRAFFIC (any server -> master -> dashboard): what a server sent and received over the last few seconds, one + * entry per second, plus the busiest message types, its HTTP routes (dashboard, UGC), RakNet's connection statistics + * and a few gauges. Sent every REPORT_SECONDS; master sends its own straight to the dashboard. + */ +struct ServerTraffic : public LUBitStream { + ServerTraffic() : LUBitStream(ServiceType::MASTER, MessageType::Master::SERVER_TRAFFIC) {} + + static constexpr int64_t REPORT_SECONDS = 5; + static constexpr uint16_t MAX_SECONDS = 180; + static constexpr uint16_t MAX_MESSAGES = 128; + static constexpr uint16_t MAX_ROUTES = 128; + static constexpr uint16_t MAX_GAUGES = 32; + static constexpr uint16_t MAX_TEXT = 200; + + ServiceType serverType{}; + uint32_t zoneId{}; + uint32_t instanceId{}; + TrafficStats::Report report; + + static void WriteHistogram(RakNet::BitStream& stream, const TrafficStats::Histogram& histogram) { + const auto sparse = histogram.Sparse(); + stream.Write(static_cast(sparse.size())); // at most BUCKETS (58) + for (const auto& [bucket, count] : sparse) { + stream.Write(bucket); + stream.Write(count); + } + stream.Write(histogram.Sum()); + } + + static bool ReadHistogram(RakNet::BitStream& stream, TrafficStats::Histogram& histogram) { + uint8_t count{}; + if (!stream.Read(count) || count > TrafficStats::Histogram::BUCKETS) return false; + std::vector> sparse(count); + for (auto& [bucket, n] : sparse) { + if (!stream.Read(bucket) || !stream.Read(n) || bucket >= TrafficStats::Histogram::BUCKETS) return false; + } + uint64_t sum{}; + if (!stream.Read(sum)) return false; + histogram = TrafficStats::Histogram::FromSparse(sparse, sum); + return true; + } + + static void WriteText(RakNet::BitStream& stream, const std::string& text) { + const auto length = static_cast(std::min(text.size(), MAX_TEXT)); + stream.Write(length); + stream.Write(text.data(), length); + } + + static bool ReadText(RakNet::BitStream& stream, std::string& text) { + uint16_t length{}; + if (!stream.Read(length) || length > MAX_TEXT) return false; + text.resize(length); + return length == 0 || stream.Read(text.data(), length); + } + + void Serialize(RakNet::BitStream& stream) const override { + stream.Write(serverType); + stream.Write(zoneId); + stream.Write(instanceId); + + const auto seconds = std::min(report.seconds.size(), MAX_SECONDS); + stream.Write(static_cast(seconds)); + // The newest seconds when there are too many + for (size_t i = report.seconds.size() - seconds; i < report.seconds.size(); i++) { + const auto& s = report.seconds[i]; + stream.Write(s.time); + stream.Write(s.packetsIn); + stream.Write(s.packetsOut); + stream.Write(s.bytesIn); + stream.Write(s.bytesOut); + stream.Write(s.httpRequests); + for (const auto status : s.httpStatus) stream.Write(status); + stream.Write(s.httpBytesOut); + WriteHistogram(stream, s.httpLatency); + } + + const auto messages = std::min(report.messages.size(), MAX_MESSAGES); + stream.Write(static_cast(messages)); + for (size_t i = 0; i < messages; i++) { + const auto& m = report.messages[i]; + stream.Write(m.key.Packed()); + stream.Write(m.count); + stream.Write(m.bytes); + } + + const auto routes = std::min(report.routes.size(), MAX_ROUTES); + stream.Write(static_cast(routes)); + for (size_t i = 0; i < routes; i++) { + const auto& r = report.routes[i]; + WriteText(stream, r.route); + stream.Write(r.count); + for (const auto status : r.status) stream.Write(status); + stream.Write(r.bytesOut); + WriteHistogram(stream, r.latency); + } + + const auto& l = report.link; + stream.Write(l.connections); + stream.Write(l.datagramsSent); + stream.Write(l.datagramsReceived); + stream.Write(l.bytesSent); + stream.Write(l.bytesReceived); + stream.Write(l.resends); + stream.Write(l.resendQueue); + stream.Write(l.averagePingMs); + + const auto gauges = std::min(report.gauges.size(), MAX_GAUGES); + stream.Write(static_cast(gauges)); + for (size_t i = 0; i < gauges; i++) { + WriteText(stream, report.gauges[i].first); + stream.Write(report.gauges[i].second); + } + } + + bool Deserialize(RakNet::BitStream& stream) override { + if (!stream.Read(serverType) || !stream.Read(zoneId) || !stream.Read(instanceId)) return false; + + uint16_t count{}; + if (!stream.Read(count) || count > MAX_SECONDS) return false; + report.seconds.resize(count); + for (auto& s : report.seconds) { + if (!stream.Read(s.time) || !stream.Read(s.packetsIn) || !stream.Read(s.packetsOut) || !stream.Read(s.bytesIn) || !stream.Read(s.bytesOut) || + !stream.Read(s.httpRequests)) return false; + for (auto& status : s.httpStatus) if (!stream.Read(status)) return false; + if (!stream.Read(s.httpBytesOut) || !ReadHistogram(stream, s.httpLatency)) return false; + } + + if (!stream.Read(count) || count > MAX_MESSAGES) return false; + report.messages.resize(count); + for (auto& m : report.messages) { + uint64_t packed{}; + if (!stream.Read(packed) || !stream.Read(m.count) || !stream.Read(m.bytes)) return false; + m.key = TrafficStats::MessageKey::Unpack(packed); + } + + if (!stream.Read(count) || count > MAX_ROUTES) return false; + report.routes.resize(count); + for (auto& r : report.routes) { + if (!ReadText(stream, r.route) || !stream.Read(r.count)) return false; + for (auto& status : r.status) if (!stream.Read(status)) return false; + if (!stream.Read(r.bytesOut) || !ReadHistogram(stream, r.latency)) return false; + } + + auto& l = report.link; + if (!stream.Read(l.connections) || !stream.Read(l.datagramsSent) || !stream.Read(l.datagramsReceived) || !stream.Read(l.bytesSent) || + !stream.Read(l.bytesReceived) || !stream.Read(l.resends) || !stream.Read(l.resendQueue) || !stream.Read(l.averagePingMs)) return false; + + if (!stream.Read(count) || count > MAX_GAUGES) return false; + report.gauges.resize(count); + for (auto& [name, value] : report.gauges) { + if (!ReadText(stream, name) || !stream.Read(value)) return false; + } + return true; + } +}; + +#endif //!__SERVERTRAFFIC__H__ diff --git a/tests/dCommonTests/CMakeLists.txt b/tests/dCommonTests/CMakeLists.txt index a08e47bcc..1e18043e1 100644 --- a/tests/dCommonTests/CMakeLists.txt +++ b/tests/dCommonTests/CMakeLists.txt @@ -28,6 +28,7 @@ set(DCOMMONTEST_SOURCES "PropertyRentRulesTests.cpp" "PropertyReputationRulesTests.cpp" "BindAddressTests.cpp" + "TrafficStatsTests.cpp" ) add_subdirectory(dEnumsTests) diff --git a/tests/dCommonTests/TrafficStatsTests.cpp b/tests/dCommonTests/TrafficStatsTests.cpp new file mode 100644 index 000000000..420cddcb2 --- /dev/null +++ b/tests/dCommonTests/TrafficStatsTests.cpp @@ -0,0 +1,228 @@ +#include + +#include +#include + +#include "TrafficStats.h" +#include "MessageIdentifiers.h" +#include "ServiceType.h" +#include "MessageType/Client.h" +#include "MessageType/Game.h" +#include "MessageType/World.h" + +using namespace TrafficStats; + +namespace { + // An LU packet header (and, for game messages, the target object and message ID) + std::vector LuPacket(ServiceType service, uint32_t packet, int32_t gameMessage = -1, size_t pad = 0) { + std::vector data{ ID_USER_PACKET_ENUM }; + const auto s = static_cast(service); + data.push_back(static_cast(s)); + data.push_back(static_cast(s >> 8)); + for (int i = 0; i < 4; i++) data.push_back(static_cast(packet >> (8 * i))); + data.push_back(0); + if (gameMessage >= 0) { + for (int i = 0; i < 8; i++) data.push_back(0x11); + data.push_back(static_cast(gameMessage)); + data.push_back(static_cast(gameMessage >> 8)); + } + data.resize(data.size() + pad); + return data; + } +} + +TEST(TrafficStatsTest, HistogramBucketsDoubleEveryThird) { + EXPECT_EQ(Histogram::UpperBound(0), 100u); + EXPECT_EQ(Histogram::UpperBound(3), 200u); + EXPECT_EQ(Histogram::UpperBound(30), 102400u); + EXPECT_EQ(Histogram::UpperBound(Histogram::BUCKETS - 1), UINT64_MAX); + EXPECT_EQ(Histogram::BucketFor(0), 0u); + EXPECT_EQ(Histogram::BucketFor(100), 0u); + EXPECT_EQ(Histogram::BucketFor(101), 1u); + EXPECT_EQ(Histogram::BucketFor(200), 3u); + EXPECT_EQ(Histogram::BucketFor(UINT64_MAX), Histogram::BUCKETS - 1); + for (size_t i = 1; i + 1 < Histogram::BUCKETS; i++) EXPECT_GT(Histogram::UpperBound(i), Histogram::UpperBound(i - 1)); +} + +TEST(TrafficStatsTest, PercentilesAreWithinABucket) { + Histogram h; + for (uint64_t ms = 1; ms <= 1000; ms++) h.Add(ms * 1000); + EXPECT_EQ(h.Count(), 1000u); + EXPECT_EQ(h.Sum(), 500500000u); + // A bucket spans 26%, so the answer is within that of the exact value + for (const auto [fraction, exact] : { std::pair{ 0.5, 500000.0 }, std::pair{ 0.95, 950000.0 }, std::pair{ 0.99, 990000.0 } }) { + const auto p = static_cast(h.Percentile(fraction)); + EXPECT_NEAR(p, exact, exact * 0.26) << fraction; + } + EXPECT_LE(h.Percentile(0.5), h.Percentile(0.95)); + EXPECT_LE(h.Percentile(0.95), h.Percentile(0.99)); + EXPECT_EQ(Histogram().Percentile(0.5), 0u); +} + +TEST(TrafficStatsTest, SinglePercentileStaysInItsBucket) { + Histogram h; + h.Add(150); + const auto p = h.Percentile(0.99); + EXPECT_GT(p, Histogram::UpperBound(0)); + EXPECT_LE(p, Histogram::UpperBound(Histogram::BucketFor(150))); + // Overflow reports its lower bound instead of infinity + Histogram slow; + slow.Add(3600ull * 1000000); + EXPECT_EQ(slow.Percentile(0.5), Histogram::UpperBound(Histogram::BUCKETS - 2)); +} + +TEST(TrafficStatsTest, HistogramsMergeAndSurviveSparse) { + Histogram a, b; + a.Add(500, 3); + b.Add(500); + b.Add(40000, 2); + a.Merge(b); + EXPECT_EQ(a.Count(), 6u); + EXPECT_EQ(a.Sum(), 500u * 4 + 80000u); + const auto sparse = a.Sparse(); + ASSERT_EQ(sparse.size(), 2u); + const auto back = Histogram::FromSparse(sparse, a.Sum()); + for (size_t i = 0; i < Histogram::BUCKETS; i++) EXPECT_EQ(back.At(i), a.At(i)); + EXPECT_EQ(back.Sum(), a.Sum()); + EXPECT_EQ(back.Percentile(0.5), a.Percentile(0.5)); +} + +TEST(TrafficStatsTest, StatusClasses) { + EXPECT_EQ(StatusClass(101), 0u); + EXPECT_EQ(StatusClass(200), 1u); + EXPECT_EQ(StatusClass(304), 2u); + EXPECT_EQ(StatusClass(404), 3u); + EXPECT_EQ(StatusClass(503), 4u); + EXPECT_EQ(StatusClass(0), 4u); + EXPECT_EQ(StatusClass(999), 4u); +} + +TEST(TrafficStatsTest, KeysFromPackets) { + const auto world = LuPacket(ServiceType::WORLD, static_cast(MessageType::World::POSITION_UPDATE), -1, 20); + auto key = KeyOf(world.data(), world.size(), false); + EXPECT_EQ(key.service, static_cast(ServiceType::WORLD)); + EXPECT_EQ(key.packet, static_cast(MessageType::World::POSITION_UPDATE)); + EXPECT_EQ(key.gameMessage, 0); + EXPECT_FALSE(key.outbound); + + const auto gm = LuPacket(ServiceType::CLIENT, static_cast(MessageType::Client::GAME_MSG), static_cast(MessageType::Game::REQUEST_USE)); + key = KeyOf(gm.data(), gm.size(), true); + EXPECT_EQ(key.service, static_cast(ServiceType::CLIENT)); + EXPECT_EQ(key.gameMessage, static_cast(MessageType::Game::REQUEST_USE)); + EXPECT_TRUE(key.outbound); + + // A game message cut short has no message ID; never reads past the end + key = KeyOf(gm.data(), 17, true); + EXPECT_EQ(key.gameMessage, 0); + + const uint8_t replica[] = { ID_REPLICA_MANAGER_SERIALIZE, 1, 2 }; + key = KeyOf(replica, sizeof(replica), true); + EXPECT_EQ(key.service, MessageKey::RAKNET); + EXPECT_EQ(key.packet, static_cast(ID_REPLICA_MANAGER_SERIALIZE)); + + const uint8_t shortLu[] = { ID_USER_PACKET_ENUM, 4, 0 }; + EXPECT_EQ(KeyOf(shortLu, sizeof(shortLu), false).service, MessageKey::RAKNET); + EXPECT_EQ(KeyOf(nullptr, 0, false).service, MessageKey::RAKNET); +} + +TEST(TrafficStatsTest, KeysPackAndUnpack) { + for (const auto& key : { MessageKey{ true, 5, 12, 1234 }, MessageKey{ false, MessageKey::RAKNET, 36, 0 }, MessageKey{ false, 4, 0xFFFFFFFF, 0xFFFF } }) { + EXPECT_EQ(MessageKey::Unpack(key.Packed()), key); + } + EXPECT_NE((MessageKey{ true, 5, 12, 0 }).Packed(), (MessageKey{ false, 5, 12, 0 }).Packed()); +} + +TEST(TrafficStatsTest, RecorderFillsSilentSecondsAndKeepsTheCurrentOne) { + Recorder r; + const MessageKey in{ false, 4, 5, 0 }; + const MessageKey out{ true, 5, 12, 0 }; + r.Packet(1000, in, 100); + r.Packet(1000, in, 50); + r.Packet(1000, out, 30, 4); // a broadcast to four + r.Packet(1003, in, 10); + r.Packet(1005, in, 10); // the current second: stays + const auto report = r.Take(1005); + ASSERT_EQ(report.seconds.size(), 5u); // 1000..1004 + EXPECT_EQ(report.seconds[0].time, 1000); + EXPECT_EQ(report.seconds[0].packetsIn, 2u); + EXPECT_EQ(report.seconds[0].bytesIn, 150u); + EXPECT_EQ(report.seconds[0].packetsOut, 4u); + EXPECT_EQ(report.seconds[0].bytesOut, 120u); + EXPECT_TRUE(report.seconds[1].Idle()); + EXPECT_EQ(report.seconds[3].packetsIn, 1u); + EXPECT_EQ(report.seconds[4].time, 1004); + + const auto next = r.Take(1008); + ASSERT_EQ(next.seconds.size(), 3u); // 1005..1007, no second twice + EXPECT_EQ(next.seconds[0].time, 1005); + EXPECT_EQ(next.seconds[0].packetsIn, 1u); +} + +TEST(TrafficStatsTest, RecorderCapsLongSilences) { + Recorder r; + r.Packet(1000, MessageKey{}, 1); + r.Take(1001); + const auto report = r.Take(1001 + 10000); + EXPECT_EQ(report.seconds.size(), static_cast(Recorder::MAX_GAP)); + EXPECT_EQ(report.seconds.back().time, 1000 + 10000); +} + +TEST(TrafficStatsTest, RecorderTopMessagesPerDirection) { + Recorder r; + for (uint32_t id = 0; id < 40; id++) { + for (uint32_t n = 0; n <= id; n++) { + r.Packet(1, MessageKey{ false, 4, id, 0 }, 10); + r.Packet(1, MessageKey{ true, 5, id, 0 }, 10); + } + } + const auto report = r.Take(2); + ASSERT_EQ(report.messages.size(), Recorder::TOP_MESSAGES * 2); + EXPECT_FALSE(report.messages.front().key.outbound); + EXPECT_EQ(report.messages.front().key.packet, 39u); + EXPECT_EQ(report.messages.front().count, 40u); + EXPECT_TRUE(report.messages[Recorder::TOP_MESSAGES].key.outbound); + EXPECT_TRUE(r.Take(3).messages.empty()); +} + +TEST(TrafficStatsTest, RecorderHttpAndRoutes) { + Recorder r; + r.Http(10, "GET /api/players", 200, 1500, 2000); + r.Http(10, "GET /api/players", 404, 300, 20); + r.Http(11, "POST /api/login", 500, 90000, 50); + for (size_t i = 0; i < Recorder::MAX_ROUTES + 5; i++) r.Http(11, "GET /r" + std::to_string(i), 200, 100, 1); + r.SetGauge("workers_busy", [] { return 3.0; }); + const auto report = r.Take(12); + ASSERT_EQ(report.seconds.size(), 2u); + EXPECT_EQ(report.seconds[0].httpRequests, 2u); + EXPECT_EQ(report.seconds[0].httpStatus[1], 1u); + EXPECT_EQ(report.seconds[0].httpStatus[3], 1u); + EXPECT_EQ(report.seconds[0].httpBytesOut, 2020u); + EXPECT_EQ(report.seconds[0].httpLatency.Count(), 2u); + EXPECT_EQ(report.routes.size(), Recorder::MAX_ROUTES + 1); // the rest counted as "other" + const auto other = std::find_if(report.routes.begin(), report.routes.end(), [](const RouteStats& s) { return s.route == "other"; }); + ASSERT_NE(other, report.routes.end()); + EXPECT_EQ(other->count, 7u); + ASSERT_EQ(report.gauges.size(), 1u); + EXPECT_EQ(report.gauges[0].second, 3.0); +} + +TEST(TrafficStatsTest, DueAfterTheInterval) { + Recorder r; + EXPECT_FALSE(r.Due(100, 5)); // starts the clock + EXPECT_FALSE(r.Due(104, 5)); + EXPECT_TRUE(r.Due(105, 5)); + r.Take(105); + EXPECT_FALSE(r.Due(106, 5)); +} + +// Counting is on every packet's path: keep it cheap +TEST(TrafficStatsTest, CountingIsCheap) { + Recorder r; + const auto packet = LuPacket(ServiceType::CLIENT, static_cast(MessageType::Client::GAME_MSG), 100, 30); + constexpr int N = 1000000; + const auto start = std::chrono::steady_clock::now(); + for (int i = 0; i < N; i++) r.Packet(Now(), KeyOf(packet.data(), packet.size(), (i & 1) != 0), packet.size()); + const auto ns = std::chrono::duration_cast(std::chrono::steady_clock::now() - start).count() / N; + std::printf("[ ] KeyOf + Recorder::Packet: %lld ns per packet\n", static_cast(ns)); + EXPECT_LT(ns, 2000); // generous for slow CI machines and sanitizers +} diff --git a/tests/dGameTests/dNetTests/CMakeLists.txt b/tests/dGameTests/dNetTests/CMakeLists.txt index 9bc202c14..514d1ad44 100644 --- a/tests/dGameTests/dNetTests/CMakeLists.txt +++ b/tests/dGameTests/dNetTests/CMakeLists.txt @@ -2,6 +2,7 @@ SET(DNET_TESTS "ChatPacketsTests.cpp" "CommonAuthPacketsTests.cpp" "MasterPacketsTests.cpp" + "ServerTrafficTests.cpp" "WorldPacketsTests.cpp") # Get the folder name and prepend it to the files above diff --git a/tests/dGameTests/dNetTests/ServerTrafficTests.cpp b/tests/dGameTests/dNetTests/ServerTrafficTests.cpp new file mode 100644 index 000000000..b7c6fd330 --- /dev/null +++ b/tests/dGameTests/dNetTests/ServerTrafficTests.cpp @@ -0,0 +1,56 @@ +#include + +#include "master/ServerTraffic.h" + +using namespace TrafficStats; + +TEST(ServerTrafficTest, ServerTrafficRoundTrips) { + ServerTraffic sent; + sent.serverType = ServiceType::WORLD; + sent.zoneId = 1100; + sent.instanceId = 3; + Second second{ .time = 1700000000, .packetsIn = 5, .packetsOut = 9, .bytesIn = 500, .bytesOut = 9000, .httpRequests = 2 }; + second.httpStatus[1] = 2; + second.httpLatency.Add(1200); + second.httpLatency.Add(30000); + sent.report.seconds = { second, Second{ .time = 1700000001 } }; + sent.report.messages = { MessageCount{ MessageKey{ true, 5, 12, 1234 }, 7, 700 } }; + RouteStats route{ .route = "GET /api/players/:id", .count = 2, .bytesOut = 99 }; + route.status[1] = 2; + route.latency.Add(800, 2); + sent.report.routes = { route }; + sent.report.link = { 12, 100, 90, 10000, 9000, 3, 1, 42 }; + sent.report.gauges = { { "workers_busy", 2.0 } }; + + RakNet::BitStream stream; + sent.WritePacket(stream); + LUBitStream header; + ASSERT_TRUE(header.ReadHeader(stream)); + EXPECT_EQ(header.internalPacketID, static_cast(MessageType::Master::SERVER_TRAFFIC)); + ServerTraffic got; + ASSERT_TRUE(got.Deserialize(stream)); + EXPECT_EQ(got.serverType, ServiceType::WORLD); + EXPECT_EQ(got.zoneId, 1100u); + EXPECT_EQ(got.instanceId, 3u); + ASSERT_EQ(got.report.seconds.size(), 2u); + EXPECT_EQ(got.report.seconds[0].bytesOut, 9000u); + EXPECT_EQ(got.report.seconds[0].httpStatus[1], 2u); + EXPECT_EQ(got.report.seconds[0].httpLatency.Count(), 2u); + EXPECT_EQ(got.report.seconds[0].httpLatency.Sum(), 31200u); + ASSERT_EQ(got.report.messages.size(), 1u); + EXPECT_EQ(got.report.messages[0].key, (MessageKey{ true, 5, 12, 1234 })); + ASSERT_EQ(got.report.routes.size(), 1u); + EXPECT_EQ(got.report.routes[0].route, "GET /api/players/:id"); + EXPECT_EQ(got.report.routes[0].latency.Count(), 2u); + EXPECT_EQ(got.report.link.averagePingMs, 42u); + EXPECT_EQ(got.report.link.connections, 12u); + ASSERT_EQ(got.report.gauges.size(), 1u); + EXPECT_EQ(got.report.gauges[0].first, "workers_busy"); + + // Truncated reports are refused + RakNet::BitStream partial(stream.GetData(), stream.GetNumberOfBytesUsed() - 3, true); + LUBitStream skip; + ASSERT_TRUE(skip.ReadHeader(partial)); + ServerTraffic broken; + EXPECT_FALSE(broken.Deserialize(partial)); +}