#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. * * Newer servers append optional sections at the end, each after a marker byte, so readers that don't know them stop * before them and reports without them still read: PEER_SPLIT_MARKER, each second's packets by peer (clients, master, * other servers) and its HTTP requests from and to other servers (peerSplit); CONNECTIONS_MARKER, the busiest remote * ends with the rest summed (hasConnections). */ 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; static constexpr uint8_t PEER_SPLIT_MARKER = 1; static constexpr uint8_t HTTP_SPLIT_BIT = 0x80; // in a second's mask: the HTTP split follows the peers static constexpr uint8_t CONNECTIONS_MARKER = 2; static constexpr uint8_t MAX_CONNECTIONS = 64; 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); } if (report.peerSplit) WritePeerSplit(stream, report.seconds.size() - seconds); if (report.hasConnections) WriteConnections(stream); } static uint32_t Clamp(uint64_t value) { return static_cast(std::min(value, UINT32_MAX)); } // Per second (from `first`, the same ones as above): a bit per peer class with traffic, then its counts void WritePeerSplit(RakNet::BitStream& stream, size_t first) const { stream.Write(PEER_SPLIT_MARKER); for (size_t i = first; i < report.seconds.size(); i++) { const auto& second = report.seconds[i]; const auto& peers = second.peers; uint8_t mask = 0; for (size_t p = 0; p < peers.size(); p++) if (!peers[p].Empty()) mask |= static_cast(1u << p); const bool http = second.httpFromServers || second.httpOutRequests; if (http) mask |= HTTP_SPLIT_BIT; stream.Write(mask); for (size_t p = 0; p < peers.size(); p++) { if (!(mask & (1u << p))) continue; stream.Write(Clamp(peers[p].packetsIn)); stream.Write(Clamp(peers[p].packetsOut)); stream.Write(Clamp(peers[p].bytesIn)); stream.Write(Clamp(peers[p].bytesOut)); } if (http) { stream.Write(Clamp(second.httpFromServers)); stream.Write(Clamp(second.httpFromServersBytesOut)); stream.Write(Clamp(second.httpOutRequests)); stream.Write(Clamp(second.httpOutBytesIn)); } } } bool ReadPeerSplit(RakNet::BitStream& stream) { for (auto& s : report.seconds) { uint8_t mask{}; if (!stream.Read(mask)) return false; for (size_t p = 0; p < s.peers.size(); p++) { if (!(mask & (1u << p))) continue; uint32_t pin{}, pout{}, bin{}, bout{}; if (!stream.Read(pin) || !stream.Read(pout) || !stream.Read(bin) || !stream.Read(bout)) return false; s.peers[p] = { pin, pout, bin, bout }; } if (mask & HTTP_SPLIT_BIT) { uint32_t from{}, fromBytes{}, out{}, outBytes{}; if (!stream.Read(from) || !stream.Read(fromBytes) || !stream.Read(out) || !stream.Read(outBytes)) return false; s.httpFromServers = from; s.httpFromServersBytesOut = fromBytes; s.httpOutRequests = out; s.httpOutBytesIn = outBytes; } } report.peerSplit = true; return true; } static void WriteConnection(RakNet::BitStream& stream, const TrafficStats::Connection& c) { stream.Write(Clamp(c.packetsIn)); stream.Write(Clamp(c.packetsOut)); stream.Write(c.bytesIn); stream.Write(c.bytesOut); stream.Write(c.resends); } static bool ReadConnection(RakNet::BitStream& stream, TrafficStats::Connection& c) { uint32_t pin{}, pout{}; if (!stream.Read(pin) || !stream.Read(pout) || !stream.Read(c.bytesIn) || !stream.Read(c.bytesOut) || !stream.Read(c.resends)) return false; c.packetsIn = pin; c.packetsOut = pout; return true; } void WriteConnections(RakNet::BitStream& stream) const { stream.Write(CONNECTIONS_MARKER); const auto count = std::min(report.connections.size(), MAX_CONNECTIONS); stream.Write(static_cast(count)); for (size_t i = 0; i < count; i++) { const auto& c = report.connections[i]; WriteText(stream, c.address); stream.Write(c.port); stream.Write(static_cast(static_cast(c.peer) | (c.http ? 0x80 : 0))); WriteConnection(stream, c); stream.Write(c.pingMs); stream.Write(c.accountId); stream.Write(c.characterId); WriteText(stream, c.account); WriteText(stream, c.character); } // The rest summed (those over MAX_CONNECTIONS too) auto others = report.otherConnections; uint32_t otherCount = report.otherConnectionCount; for (size_t i = count; i < report.connections.size(); i++, otherCount++) others.Merge(report.connections[i]); stream.Write(otherCount); WriteConnection(stream, others); } bool ReadConnections(RakNet::BitStream& stream) { uint8_t count{}; if (!stream.Read(count) || count > MAX_CONNECTIONS) return false; report.connections.resize(count); for (auto& c : report.connections) { uint8_t flags{}; if (!ReadText(stream, c.address) || !stream.Read(c.port) || !stream.Read(flags) || !ReadConnection(stream, c) || !stream.Read(c.pingMs) || !stream.Read(c.accountId) || !stream.Read(c.characterId) || !ReadText(stream, c.account) || !ReadText(stream, c.character)) return false; c.peer = static_cast(std::min(flags & 0x7F, TrafficStats::PEER_CLASSES - 1)); c.http = (flags & 0x80) != 0; } if (!stream.Read(report.otherConnectionCount) || !ReadConnection(stream, report.otherConnections)) return false; report.hasConnections = true; return true; } 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; } // Older servers stop here; newer ones add sections, each after its marker (a reader stops at one it doesn't know) report.peerSplit = false; report.hasConnections = false; uint8_t marker{}; while (stream.GetNumberOfUnreadBits() >= 8 && stream.Read(marker)) { if (marker == PEER_SPLIT_MARKER && !report.peerSplit) { if (!ReadPeerSplit(stream)) return false; } else if (marker == CONNECTIONS_MARKER && !report.hasConnections) { if (!ReadConnections(stream)) return false; } else { break; } } return true; } }; #endif //!__SERVERTRAFFIC__H__