feat(traffic): send the peer split and connections in SERVER_TRAFFIC

Two optional sections at the end of the report, each after a marker
byte: every second's packets by peer with its HTTP requests from and to
other servers, and the busiest remote ends with the rest summed.
Reports without them still read (older servers), and older readers stop
before them.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Aaron Kimbrell
2026-09-29 18:23:31 -05:00
parent 9e8f161dd8
commit 3b778065cf
2 changed files with 197 additions and 0 deletions

View File

@@ -16,6 +16,11 @@
* 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) {}
@@ -26,6 +31,10 @@ struct ServerTraffic : public LUBitStream {
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{};
@@ -125,6 +134,117 @@ struct ServerTraffic : public LUBitStream {
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<uint32_t>(std::min<uint64_t>(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<uint8_t>(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<size_t>(report.connections.size(), MAX_CONNECTIONS);
stream.Write(static_cast<uint8_t>(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<uint8_t>(static_cast<uint8_t>(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<TrafficStats::Peer>(std::min<uint8_t>(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 {
@@ -165,6 +285,20 @@ struct ServerTraffic : public LUBitStream {
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;
}
};

View File

@@ -1,4 +1,5 @@
#include <gtest/gtest.h>
#include <cstring>
#include "master/ServerTraffic.h"
@@ -54,3 +55,65 @@ TEST(ServerTrafficTest, ServerTrafficRoundTrips) {
ServerTraffic broken;
EXPECT_FALSE(broken.Deserialize(partial));
}
namespace {
ServerTraffic Sample(bool split) {
ServerTraffic sent;
sent.serverType = ServiceType::AUTH;
Second a{ .time = 1700000000, .packetsIn = 4, .packetsOut = 3, .bytesIn = 400, .bytesOut = 300 };
a.peers[0] = { 3, 2, 300, 200 };
a.peers[1] = { 1, 1, 100, 100 };
a.httpFromServers = 2;
a.httpFromServersBytesOut = 64;
a.httpOutRequests = 1;
a.httpOutBytesIn = 9;
sent.report.seconds = { a, Second{ .time = 1700000001 } };
sent.report.gauges = { { "workers_busy", 1.0 } };
sent.report.peerSplit = split;
return sent;
}
bool ReadBack(RakNet::BitStream& stream, size_t bytes, ServerTraffic& got) {
RakNet::BitStream in(stream.GetData(), bytes, true);
LUBitStream header;
return header.ReadHeader(in) && got.Deserialize(in);
}
}
TEST(ServerTrafficTest, PeerSplitRoundTrips) {
RakNet::BitStream stream;
Sample(true).WritePacket(stream);
ServerTraffic got;
ASSERT_TRUE(ReadBack(stream, stream.GetNumberOfBytesUsed(), got));
EXPECT_TRUE(got.report.peerSplit);
ASSERT_EQ(got.report.seconds.size(), 2u);
EXPECT_EQ(got.report.seconds[0].peers[0], (PeerCounts{ 3, 2, 300, 200 }));
EXPECT_EQ(got.report.seconds[0].peers[1], (PeerCounts{ 1, 1, 100, 100 }));
EXPECT_TRUE(got.report.seconds[0].peers[2].Empty());
EXPECT_EQ(got.report.seconds[0].httpFromServers, 2u);
EXPECT_EQ(got.report.seconds[0].httpFromServersBytesOut, 64u);
EXPECT_EQ(got.report.seconds[0].httpOutRequests, 1u);
EXPECT_EQ(got.report.seconds[0].httpOutBytesIn, 9u);
EXPECT_TRUE(got.report.seconds[1].peers[0].Empty());
EXPECT_EQ(got.report.gauges.size(), 1u);
// A cut-off split is refused
ServerTraffic broken;
EXPECT_FALSE(ReadBack(stream, stream.GetNumberOfBytesUsed() - 2, broken));
}
TEST(ServerTrafficTest, ReportsWithoutTheSplitStillRead) {
// An older server's report is the same bytes without the end: it reads, with no split
RakNet::BitStream old, now;
Sample(false).WritePacket(old);
Sample(true).WritePacket(now);
ASSERT_LT(old.GetNumberOfBytesUsed(), now.GetNumberOfBytesUsed());
EXPECT_EQ(0, std::memcmp(old.GetData(), now.GetData(), old.GetNumberOfBytesUsed()));
ServerTraffic got;
ASSERT_TRUE(ReadBack(old, old.GetNumberOfBytesUsed(), got));
EXPECT_FALSE(got.report.peerSplit);
EXPECT_EQ(got.report.seconds[0].packetsIn, 4u);
EXPECT_TRUE(got.report.seconds[0].peers[0].Empty());
EXPECT_EQ(got.report.gauges.size(), 1u);
}