feat(net): count every server's packets and send a traffic report every 5 seconds

TrafficStats keeps packets and bytes in and out per second and the busiest
packet and game message types. dServer counts at its send and receive calls
and adds RakNet's connection statistics (datagrams, resends, ping); the report
goes to master as SERVER_TRAFFIC (appended), which passes it to the dashboard.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Aaron Kimbrell
2026-09-27 08:44:03 -05:00
parent 623ab58afb
commit 13ad679df1
12 changed files with 1015 additions and 1 deletions

276
dCommon/TrafficStats.cpp Normal file
View File

@@ -0,0 +1,276 @@
#include "TrafficStats.h"
#include <algorithm>
#include <cmath>
#include <ctime>
#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<uint64_t, Histogram::BUCKETS>& Bounds() {
static const auto bounds = [] {
std::array<uint64_t, Histogram::BUCKETS> out{};
for (size_t i = 0; i + 1 < Histogram::BUCKETS; i++) out[i] = static_cast<uint64_t>(std::llround(FIRST_BOUND * std::exp2(static_cast<double>(i) / 3.0)));
out[Histogram::BUCKETS - 1] = UINT64_MAX;
return out;
}();
return bounds;
}
template<size_t N>
void AddArrays(std::array<uint64_t, N>& to, const std::array<uint64_t, N>& 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<size_t>(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<double>(m_Count)));
uint64_t seen = 0;
for (size_t i = 0; i < BUCKETS; i++) {
if (m_Counts[i] == 0) continue;
if (static_cast<double>(seen + m_Counts[i]) >= rank) {
const double lower = i == 0 ? 0.0 : static_cast<double>(UpperBound(i - 1));
// The overflow bucket has no upper bound: report its lower one
if (i == BUCKETS - 1) return static_cast<uint64_t>(lower);
const double upper = static_cast<double>(UpperBound(i));
const double within = (rank - static_cast<double>(seen)) / static_cast<double>(m_Counts[i]);
return static_cast<uint64_t>(std::llround(lower + (upper - lower) * within));
}
seen += m_Counts[i];
}
return 0;
}
std::vector<std::pair<uint8_t, uint32_t>> Histogram::Sparse() const {
std::vector<std::pair<uint8_t, uint32_t>> out;
for (size_t i = 0; i < BUCKETS; i++) {
if (m_Counts[i]) out.emplace_back(static_cast<uint8_t>(i), m_Counts[i]);
}
return out;
}
Histogram Histogram::FromSparse(const std::vector<std::pair<uint8_t, uint32_t>>& 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<uint64_t>(outbound) << 63) | (static_cast<uint64_t>(service & 0x7FFF) << 48) |
(static_cast<uint64_t>(packet) << 16) | gameMessage;
}
MessageKey MessageKey::Unpack(uint64_t packed) {
MessageKey key;
key.outbound = (packed >> 63) != 0;
key.service = static_cast<uint16_t>((packed >> 48) & 0x7FFF);
if (key.service == (RAKNET & 0x7FFF)) key.service = RAKNET;
key.packet = static_cast<uint32_t>(packed >> 16);
key.gameMessage = static_cast<uint16_t>(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<uint16_t>(data[1] | (data[2] << 8));
key.packet = static_cast<uint32_t>(data[3]) | (static_cast<uint32_t>(data[4]) << 8) | (static_cast<uint32_t>(data[5]) << 16) | (static_cast<uint32_t>(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<uint16_t>(ServiceType::WORLD) && key.packet == static_cast<uint32_t>(MessageType::World::GAME_MSG)) ||
(key.service == static_cast<uint16_t>(ServiceType::CLIENT) && key.packet == static_cast<uint32_t>(MessageType::Client::GAME_MSG));
if (gameMessage && length >= 18) key.gameMessage = static_cast<uint16_t>(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<double()> 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<std::pair<std::string, std::function<double()>>> 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<MessageCount> 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<MessageCount> Top(const std::vector<MessageCount>& counts, size_t limit) {
std::vector<MessageCount> out;
for (const bool outbound : { false, true }) {
std::vector<MessageCount> 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<int64_t>(std::time(nullptr));
}
}