From 9e8f161dd825b6879a179b4f1e8f7349b7517ad6 Mon Sep 17 00:00:00 2001 From: Aaron Kimbrell Date: Tue, 29 Sep 2026 18:23:31 -0500 Subject: [PATCH] feat(traffic): count traffic by peer and by connection Every server now splits its packet counts by peer: its own connections (players on auth and worlds), its master link, or other servers (the worlds on chat, every server on master, a world's chat link, which is now counted too). HTTP requests carrying X-Darkflame-Server count as another server's, the dashboard counts its own requests to the UGC server, and each HTTP client address is counted. The report also gets each RakNet connection's statistics (worlds name the player on it), trimmed to the 32 busiest with the rest summed. Counting stays on the main loop, except the dashboard's UGC fetches, which only touch the locked recorder. Co-Authored-By: Claude Opus 5.5 --- dCommon/TrafficStats.cpp | 81 +++++++++++++++++++++++- dCommon/TrafficStats.h | 71 ++++++++++++++++++++- dDashboardServer/routes/UgcFetch.cpp | 7 ++ dGame/dUtilities/ChatServerLink.cpp | 7 ++ dGame/dUtilities/ChatServerLink.h | 3 + dNet/dServer.cpp | 39 +++++++++--- dNet/dServer.h | 12 +++- dWeb/Web.cpp | 22 +++++-- dWorldServer/WorldServer.cpp | 17 +++-- tests/dCommonTests/TrafficStatsTests.cpp | 41 ++++++++++++ 10 files changed, 273 insertions(+), 27 deletions(-) diff --git a/dCommon/TrafficStats.cpp b/dCommon/TrafficStats.cpp index fd79da820..179b51bdf 100644 --- a/dCommon/TrafficStats.cpp +++ b/dCommon/TrafficStats.cpp @@ -3,6 +3,7 @@ #include #include #include +#include #include "MessageIdentifiers.h" #include "ServiceType.h" @@ -134,7 +135,15 @@ namespace TrafficStats { return key; } + void PeerCounts::Merge(const PeerCounts& other) { + packetsIn += other.packetsIn; + packetsOut += other.packetsOut; + bytesIn += other.bytesIn; + bytesOut += other.bytesOut; + } + void Second::Merge(const Second& other) { + for (size_t i = 0; i < PEER_CLASSES; i++) peers[i].Merge(other.peers[i]); packetsIn += other.packetsIn; packetsOut += other.packetsOut; bytesIn += other.bytesIn; @@ -143,6 +152,10 @@ namespace TrafficStats { AddArrays(httpStatus, other.httpStatus); httpBytesOut += other.httpBytesOut; httpLatency.Merge(other.httpLatency); + httpFromServers += other.httpFromServers; + httpFromServersBytesOut += other.httpFromServersBytesOut; + httpOutRequests += other.httpOutRequests; + httpOutBytesIn += other.httpOutBytesIn; } void RouteStats::Merge(const RouteStats& other) { @@ -152,16 +165,21 @@ namespace TrafficStats { latency.Merge(other.latency); } - void Recorder::Packet(int64_t now, const MessageKey& key, uint64_t bytes, uint32_t fanout) { + void Recorder::Packet(int64_t now, const MessageKey& key, uint64_t bytes, uint32_t fanout, Peer peer) { if (fanout == 0) return; std::lock_guard lock(m_Mutex); auto& second = SecondAt(now); + auto& side = second.peers[std::min(static_cast(peer), PEER_CLASSES - 1)]; if (key.outbound) { second.packetsOut += fanout; second.bytesOut += bytes * fanout; + side.packetsOut += fanout; + side.bytesOut += bytes * fanout; } else { second.packetsIn += fanout; second.bytesIn += bytes * fanout; + side.packetsIn += fanout; + side.bytesIn += bytes * fanout; } auto& message = m_Messages[key.Packed()]; message.key = key; @@ -169,9 +187,58 @@ namespace TrafficStats { message.bytes += bytes * fanout; } - void Recorder::Http(int64_t now, const std::string& route, uint16_t status, uint64_t microseconds, uint64_t bytesOut) { + void Connection::Merge(const Connection& other) { + packetsIn += other.packetsIn; + packetsOut += other.packetsOut; + bytesIn += other.bytesIn; + bytesOut += other.bytesOut; + resends += other.resends; + } + + void TrimConnections(Report& report, size_t limit) { + report.hasConnections = true; + auto& list = report.connections; + std::sort(list.begin(), list.end(), [](const Connection& a, const Connection& b) { + return a.Bytes() != b.Bytes() ? a.Bytes() > b.Bytes() : std::tie(a.address, a.port) < std::tie(b.address, b.port); + }); + for (size_t i = limit; i < list.size(); i++) { + report.otherConnections.Merge(list[i]); + report.otherConnectionCount++; + } + if (list.size() > limit) list.resize(limit); + } + + void Recorder::HttpClient(const std::string& address, bool fromServer, uint64_t bytesIn, uint64_t bytesOut) { + std::lock_guard lock(m_Mutex); + auto it = m_HttpClients.find(address); + if (it == m_HttpClients.end()) { + const bool full = m_HttpClients.size() >= MAX_HTTP_CLIENTS; + it = m_HttpClients.try_emplace(full ? std::string() : address).first; + it->second.address = it->first; + it->second.http = true; + } + auto& client = it->second; + if (fromServer) client.peer = Peer::SERVERS; + client.packetsIn++; + client.packetsOut++; + client.bytesIn += bytesIn; + client.bytesOut += bytesOut; + } + + void Recorder::HttpOut(int64_t now, uint64_t bytesIn) { std::lock_guard lock(m_Mutex); auto& second = SecondAt(now); + second.httpOutRequests++; + second.httpOutBytesIn += bytesIn; + } + + void Recorder::Http(int64_t now, const std::string& route, uint16_t status, uint64_t microseconds, uint64_t bytesOut, bool fromServer) { + std::lock_guard lock(m_Mutex); + auto& second = SecondAt(now); + if (fromServer) { + second.httpFromServers++; + second.httpFromServersBytesOut += bytesOut; + } second.httpRequests++; second.httpStatus[StatusClass(status)]++; second.httpBytesOut += bytesOut; @@ -221,6 +288,7 @@ namespace TrafficStats { Report Recorder::Take(int64_t now) { Report report; + report.peerSplit = true; std::vector>> gauges; { std::lock_guard lock(m_Mutex); @@ -245,6 +313,15 @@ namespace TrafficStats { for (auto& [_, route] : m_Routes) report.routes.push_back(std::move(route)); m_Routes.clear(); + for (auto& [address, client] : m_HttpClients) { + if (address.empty()) { // the ones over MAX_HTTP_CLIENTS + report.otherConnections.Merge(client); + report.otherConnectionCount++; + } else { + report.connections.push_back(std::move(client)); + } + } + m_HttpClients.clear(); gauges = m_Gauges; } for (const auto& [name, source] : gauges) report.gauges.emplace_back(name, source ? source() : 0.0); diff --git a/dCommon/TrafficStats.h b/dCommon/TrafficStats.h index 36e9dc560..99ce3a197 100644 --- a/dCommon/TrafficStats.h +++ b/dCommon/TrafficStats.h @@ -55,6 +55,9 @@ namespace TrafficStats { uint64_t m_Sum{}; }; + // The header a server puts on HTTP requests it makes to another server's web server, naming itself + constexpr const char* SERVER_HEADER = "X-Darkflame-Server"; + // HTTP status classes 1xx..5xx as index 0..4 (anything else counts as 5xx) size_t StatusClass(uint16_t status); @@ -75,6 +78,25 @@ namespace TrafficStats { // Reads the key from a packet's first bytes (never reads past `length`) MessageKey KeyOf(const uint8_t* data, size_t length, bool outbound); + /** + * Which side of a server a packet went to or came from, so the dashboard can draw each link: its own connections + * (players' game clients on auth and worlds), its link to master, or other servers (every server on master, the + * worlds on chat, chat on a world). + */ + enum class Peer : uint8_t { CLIENTS = 0, MASTER = 1, SERVERS = 2 }; + constexpr size_t PEER_CLASSES = 3; + + struct PeerCounts { + uint64_t packetsIn{}; + uint64_t packetsOut{}; + uint64_t bytesIn{}; + uint64_t bytesOut{}; + + void Merge(const PeerCounts& other); + bool Empty() const { return packetsIn == 0 && packetsOut == 0; } + bool operator==(const PeerCounts&) const = default; + }; + struct MessageCount { MessageKey key; uint64_t count{}; @@ -92,6 +114,11 @@ namespace TrafficStats { std::array httpStatus{}; // by StatusClass uint64_t httpBytesOut{}; Histogram httpLatency; + std::array peers{}; // the packets above split by Peer + uint64_t httpFromServers{}; // of httpRequests, those another server made (X-Darkflame-Server) + uint64_t httpFromServersBytesOut{}; // their response bodies + uint64_t httpOutRequests{}; // HTTP requests this server made to another server's web server + uint64_t httpOutBytesIn{}; // the bodies those answered with void Merge(const Second& other); // adds the counts (keeps this one's time) bool Idle() const { return packetsIn == 0 && packetsOut == 0 && httpRequests == 0; } @@ -119,23 +146,60 @@ namespace TrafficStats { uint32_t averagePingMs{}; // over the connections }; + /** + * One remote end of a server over a report: a RakNet connection (RakNet's own counts: datagrams, with + * acknowledgements and resends) or an HTTP client address (requests and body bytes). Personal data: the dashboard + * keeps these in memory only and shows addresses only with network_ips. + */ + struct Connection { + std::string address; // "203.0.113.5" or an IPv6 address; "" for the sum of the ones left out + uint16_t port{}; // 0 for HTTP clients (one entry per address) + Peer peer{}; + bool http{}; + uint64_t packetsIn{}, packetsOut{}; // datagrams, or HTTP requests and responses + uint64_t bytesIn{}, bytesOut{}; + uint32_t resends{}; + uint32_t pingMs{}; + uint32_t accountId{}; // a logged-in player's (world servers), else 0 + uint64_t characterId{}; + std::string account, character; + + uint64_t Bytes() const { return bytesIn + bytesOut; } + void Merge(const Connection& other); // adds the counts + }; + 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 + bool peerSplit{}; // whether the seconds' `peers` are filled (reports of older servers don't have them) + std::vector connections; // the busiest remote ends, busiest first + Connection otherConnections; // the rest summed + uint32_t otherConnectionCount{}; + bool hasConnections{}; // false in reports of older servers }; + // Keeps the `limit` busiest connections of the report, sums the rest into otherConnections + void TrimConnections(Report& report, size_t limit); + 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 + static constexpr size_t MAX_HTTP_CLIENTS = 1024; // addresses per report; more count as one + static constexpr size_t TOP_CONNECTIONS = 32; // connections a report names; the rest are summed - // `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); + // `fanout`: how many connections a broadcast went to; `peer`: which side it went to or came from + void Packet(int64_t now, const MessageKey& key, uint64_t bytes, uint32_t fanout = 1, Peer peer = Peer::CLIENTS); + // `fromServer`: the request came from another server (it sent SERVER_HEADER), not a browser or game client + void Http(int64_t now, const std::string& route, uint16_t status, uint64_t microseconds, uint64_t bytesOut, bool fromServer = false); + // A request this server made to another server's web server (the dashboard to the UGC server); any thread + void HttpOut(int64_t now, uint64_t bytesIn); + // An HTTP client's request by its address (bytes of the request and of the answer's body) + void HttpClient(const std::string& address, bool fromServer, uint64_t bytesIn, uint64_t bytesOut); // Evaluated when a report is taken (on the thread that takes it) void SetGauge(const std::string& name, std::function source); @@ -159,6 +223,7 @@ namespace TrafficStats { int64_t m_LastReported{}; // last second a report covered std::unordered_map m_Messages; std::map m_Routes; + std::unordered_map m_HttpClients; std::vector>> m_Gauges; int64_t m_LastTake{}; }; diff --git a/dDashboardServer/routes/UgcFetch.cpp b/dDashboardServer/routes/UgcFetch.cpp index 7de5475f9..824bdd9cc 100644 --- a/dDashboardServer/routes/UgcFetch.cpp +++ b/dDashboardServer/routes/UgcFetch.cpp @@ -8,6 +8,7 @@ #include "dConfig.h" #include "Game.h" #include "RouteUtils.h" +#include "TrafficStats.h" #include "TtlCache.h" namespace { @@ -42,11 +43,16 @@ namespace { curl_easy_setopt(curl, CURLOPT_USERAGENT, "DarkflameServer-Dashboard"); curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, Collect); curl_easy_setopt(curl, CURLOPT_WRITEDATA, &out->body); + // Names the dashboard to the UGC server, whose traffic report then counts this as the dashboard's (not a player's) + curl_slist* identify = curl_slist_append(nullptr, (std::string(TrafficStats::SERVER_HEADER) + ": dashboard").c_str()); + curl_easy_setopt(curl, CURLOPT_HTTPHEADER, identify); setup(curl); const auto code = curl_easy_perform(curl); + TrafficStats::Local().HttpOut(TrafficStats::Now(), out->body.size()); if (code == CURLE_OK) curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &out->status); else out->error = curl_easy_strerror(code); curl_easy_cleanup(curl); + curl_slist_free_all(identify); return out; } } @@ -79,6 +85,7 @@ namespace UgcFetch { std::shared_ptr AdminPost(const std::string& url, const std::string& key, const std::string& body) { curl_slist* headers = curl_slist_append(nullptr, "Content-Type: application/json"); + headers = curl_slist_append(headers, (std::string(TrafficStats::SERVER_HEADER) + ": dashboard").c_str()); headers = curl_slist_append(headers, ("X-Ugc-Admin-Key: " + key).c_str()); auto out = Perform(url, 120L, [&](CURL* curl) { curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers); diff --git a/dGame/dUtilities/ChatServerLink.cpp b/dGame/dUtilities/ChatServerLink.cpp index 2cd56248d..6f295d525 100644 --- a/dGame/dUtilities/ChatServerLink.cpp +++ b/dGame/dUtilities/ChatServerLink.cpp @@ -4,10 +4,17 @@ #include "BitStreamUtils.h" #include "Game.h" #include "RakPeerInterface.h" +#include "TrafficStats.h" void ChatServerLink::Send(const LUBitStream& msg, const PacketPriority priority, const PacketReliability reliability) { if (!Game::chatServer) return; RakNet::BitStream bitStream; msg.WritePacket(bitStream); + const auto bytes = bitStream.GetNumberOfBytesUsed(); + TrafficStats::Local().Packet(TrafficStats::Now(), TrafficStats::KeyOf(bitStream.GetData(), bytes, true), bytes, 1, TrafficStats::Peer::SERVERS); Game::chatServer->Send(&bitStream, priority, reliability, 0, Game::chatSysAddr, false); } + +void ChatServerLink::CountReceived(const unsigned char* data, unsigned int length) { + TrafficStats::Local().Packet(TrafficStats::Now(), TrafficStats::KeyOf(data, length, false), length, 1, TrafficStats::Peer::SERVERS); +} diff --git a/dGame/dUtilities/ChatServerLink.h b/dGame/dUtilities/ChatServerLink.h index df47bc980..572a97ead 100644 --- a/dGame/dUtilities/ChatServerLink.h +++ b/dGame/dUtilities/ChatServerLink.h @@ -9,6 +9,9 @@ struct LUBitStream; namespace ChatServerLink { // World -> chat: sends msg (a ChatPackets struct) over the world's chat connection void Send(const LUBitStream& msg, PacketPriority priority = SYSTEM_PRIORITY, PacketReliability reliability = RELIABLE); + + // Counts a packet the world received from chat in its traffic diagnostics (TrafficStats, as another server's) + void CountReceived(const unsigned char* data, unsigned int length); } #endif // CHATSERVERLINK_H diff --git a/dNet/dServer.cpp b/dNet/dServer.cpp index 5f7081b61..d43ce3d1d 100644 --- a/dNet/dServer.cpp +++ b/dNet/dServer.cpp @@ -137,7 +137,7 @@ Packet* dServer::ReceiveFromMaster() { if (!mMasterConnectionActive) ConnectToMaster(); Packet* packet = mMasterPeer->Receive(); - CountTraffic(packet); + CountTraffic(packet, TrafficStats::Peer::MASTER); if (packet) { if (packet->length < 1) { mMasterPeer->DeallocatePacket(packet); return nullptr; } @@ -202,7 +202,7 @@ Packet* dServer::ReceiveFromMaster() { Packet* dServer::Receive() { Packet* packet = mPeer->Receive(); - CountTraffic(packet); + CountTraffic(packet, PeerOfConnections()); return packet; } @@ -216,13 +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); + CountTraffic(bitStream, broadcast, sysAddr, PeerOfConnections()); mPeer->Send(&bitStream, SYSTEM_PRIORITY, RELIABLE_ORDERED, 0, sysAddr, broadcast); } void dServer::SendToMaster(RakNet::BitStream& bitStream) { if (!mMasterConnectionActive) ConnectToMaster(); - CountTraffic(bitStream, false, mMasterSystemAddress); + CountTraffic(bitStream, false, mMasterSystemAddress, TrafficStats::Peer::MASTER); mMasterPeer->Send(&bitStream, SYSTEM_PRIORITY, RELIABLE_ORDERED, 0, mMasterSystemAddress, false); } @@ -231,7 +231,7 @@ void dServer::Disconnect(const SystemAddress& sysAddr, eServerDisconnectIdentifi notify.disconnectID = disconNotifyID; RakNet::BitStream bitStream; notify.WritePacket(bitStream); - CountTraffic(bitStream, false, sysAddr); + CountTraffic(bitStream, false, sysAddr, PeerOfConnections()); mPeer->Send(&bitStream, SYSTEM_PRIORITY, RELIABLE_ORDERED, 0, sysAddr, false); mPeer->CloseConnection(sysAddr, true); @@ -326,13 +326,17 @@ int dServer::GetLatestPing(const SystemAddress& sysAddr) const { return mPeer->GetLastPing(sysAddr); } -void dServer::CountTraffic(const Packet* packet) { +TrafficStats::Peer dServer::PeerOfConnections() const { + return mIsInternal || mServerType == ServiceType::CHAT ? TrafficStats::Peer::SERVERS : TrafficStats::Peer::CLIENTS; +} + +void dServer::CountTraffic(const Packet* packet, TrafficStats::Peer peer) { const auto now = TrafficStats::Now(); - if (packet) TrafficStats::Local().Packet(now, TrafficStats::KeyOf(packet->data, packet->length, false), packet->length); + if (packet) TrafficStats::Local().Packet(now, TrafficStats::KeyOf(packet->data, packet->length, false), packet->length, 1, peer); if (TrafficStats::Local().Due(now, ServerTraffic::REPORT_SECONDS)) ReportTraffic(); } -void dServer::CountTraffic(const RakNet::BitStream& bitStream, bool broadcast, const SystemAddress& sysAddr) { +void dServer::CountTraffic(const RakNet::BitStream& bitStream, bool broadcast, const SystemAddress& sysAddr, TrafficStats::Peer peer) { // A broadcast goes to every connection but the one given uint32_t fanout = 1; if (broadcast) { @@ -340,7 +344,7 @@ void dServer::CountTraffic(const RakNet::BitStream& bitStream, bool broadcast, c 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); + TrafficStats::Local().Packet(TrafficStats::Now(), TrafficStats::KeyOf(bitStream.GetData(), bytes, true), bytes, fanout, peer); } void dServer::AddLinkStats(RakPeerInterface* peer, uint64_t peerIndex, ServerTraffic& report, uint64_t& pingSum, std::map& seen) { @@ -365,8 +369,22 @@ void dServer::AddLinkStats(RakPeerInterface* peer, uint64_t peerIndex, ServerTra link.resends += delta(now.resends, before.resends); link.resendQueue += stats->messagesOnResendQueue; link.connections++; - pingSum += std::max(0, peer->GetAveragePing(addresses[i])); + const int ping = std::max(0, peer->GetAveragePing(addresses[i])); + pingSum += ping; seen[key] = now; + + TrafficStats::Connection connection; + connection.address = addresses[i].ToString(false); + connection.port = addresses[i].port; + connection.peer = peerIndex == 1 ? TrafficStats::Peer::MASTER : PeerOfConnections(); + connection.packetsIn = delta(now.datagramsReceived, before.datagramsReceived); + connection.packetsOut = delta(now.datagramsSent, before.datagramsSent); + connection.bytesIn = delta(now.bitsReceived, before.bitsReceived) / 8; + connection.bytesOut = delta(now.bitsSent, before.bitsSent) / 8; + connection.resends = static_cast(delta(now.resends, before.resends)); + connection.pingMs = static_cast(ping); + if (peerIndex == 0 && mConnectionIdentity) mConnectionIdentity(addresses[i], connection); + report.report.connections.push_back(std::move(connection)); } } @@ -382,6 +400,7 @@ void dServer::ReportTraffic() { AddLinkStats(mPeer, 0, report, pingSum, seen); AddLinkStats(mMasterPeer, 1, report, pingSum, seen); mLinkCounters = std::move(seen); + TrafficStats::TrimConnections(report.report, TrafficStats::Recorder::TOP_CONNECTIONS); if (report.report.link.connections) report.report.link.averagePingMs = static_cast(pingSum / report.report.link.connections); if (mTrafficSink) mTrafficSink(report); diff --git a/dNet/dServer.h b/dNet/dServer.h index 092dddab4..a1005d895 100644 --- a/dNet/dServer.h +++ b/dNet/dServer.h @@ -7,6 +7,7 @@ #include "RakPeerInterface.h" #include "ReplicaManager.h" #include "NetworkIDManager.h" +#include "TrafficStats.h" class Logger; class dConfig; @@ -58,6 +59,10 @@ public: using TrafficSink = std::function; void SetTrafficSink(TrafficSink sink) { mTrafficSink = std::move(sink); } + // Names who is on a connection in the traffic report (a world fills in the player's account and character) + using ConnectionIdentity = std::function; + void SetConnectionIdentity(ConnectionIdentity identity) { mConnectionIdentity = std::move(identity); } + bool IsConnected(const SystemAddress& sysAddr); const std::string& GetIP() const { return mIP; } const int GetPort() const { return mPort; } @@ -92,8 +97,10 @@ private: }; 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 CountTraffic(const Packet* packet, TrafficStats::Peer peer); + void CountTraffic(const RakNet::BitStream& bitStream, bool broadcast, const SystemAddress& sysAddr, TrafficStats::Peer peer); + // Who mPeer's connections are: other servers on master and chat (the worlds connect to chat), players elsewhere + TrafficStats::Peer PeerOfConnections() const; 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); @@ -136,6 +143,7 @@ protected: SendObserver mSendObserver; TrafficSink mTrafficSink; + ConnectionIdentity mConnectionIdentity; // RakNet's per-connection statistics are totals since the connection opened; the last ones seen, for deltas std::map mLinkCounters; }; diff --git a/dWeb/Web.cpp b/dWeb/Web.cpp index a469202de..96b04a269 100644 --- a/dWeb/Web.cpp +++ b/dWeb/Web.cpp @@ -260,6 +260,9 @@ namespace { struct DeferredTiming { std::string route; TrafficClock::time_point started; + bool fromServer{}; + std::string address; + uint64_t requestBytes{}; }; std::unordered_map g_DeferredTiming; @@ -280,9 +283,11 @@ namespace { std::filesystem::remove(reply.file, ec); } - void CountRequest(const std::string& route, uint16_t status, TrafficClock::time_point started, uint64_t bytes) { + void CountRequest(const std::string& route, uint16_t status, TrafficClock::time_point started, uint64_t bytes, bool fromServer, + const std::string& address, uint64_t requestBytes) { const auto micros = std::chrono::duration_cast(TrafficClock::now() - started).count(); - TrafficStats::Local().Http(TrafficStats::Now(), route, status, static_cast(std::max(micros, 0)), bytes); + TrafficStats::Local().Http(TrafficStats::Now(), route, status, static_cast(std::max(micros, 0)), bytes, fromServer); + TrafficStats::Local().HttpClient(address, fromServer, requestBytes, bytes); } } @@ -293,6 +298,10 @@ void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_ms const auto started = TrafficClock::now(); // The route's pattern, not the path, so traffic diagnostics have one entry per route std::string trafficRoute = "(no route)"; + // Another server asking (the dashboard fetching from the UGC server), for the network diagram's links + const bool fromServer = http_msg && mg_http_get_header(const_cast(http_msg), TrafficStats::SERVER_HEADER) != nullptr; + const auto clientAddress = GetClientIP(connection); + const uint64_t requestBytes = http_msg ? http_msg->message.len : 0; if (!http_msg) { reply.status = eHTTPStatusCode::BAD_REQUEST; @@ -379,7 +388,7 @@ void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_ms } } - CountRequest("GET /ws", level ? 101 : 401, started, 0); + CountRequest("GET /ws", level ? 101 : 401, started, 0, fromServer, clientAddress, requestBytes); if (level) { mg_ws_upgrade(connection, const_cast(http_msg), NULL); g_AuthenticatedWSConnections[connection] = { level->level, level->accountId, connectToken, apiToken, @@ -511,14 +520,14 @@ void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_ms g_Deferred.SetReplyOptions(connection->id, reply.headers, cc && mg_strcasecmp(*cc, mg_str("close")) == 0); // Requests the answers never came for (the client left) are forgotten now and then if (g_DeferredTiming.size() > 10000) g_DeferredTiming.clear(); - g_DeferredTiming[connection->id] = { std::move(trafficRoute), started }; + g_DeferredTiming[connection->id] = { std::move(trafficRoute), started, fromServer, clientAddress, requestBytes }; return; } // The handler deferred and then failed: its late answer is dropped if (g_Deferred.IsPending(connection->id)) g_Deferred.Close(connection->id); SendReply(connection, reply, http_msg); - CountRequest(trafficRoute, static_cast(reply.status), started, ReplyBytes(reply)); + CountRequest(trafficRoute, static_cast(reply.status), started, ReplyBytes(reply), fromServer, clientAddress, requestBytes); RemoveSentFile(reply); } @@ -799,7 +808,8 @@ void Web::SendDeferredReplies() { // Clears is_resp once the reply is out, so mongoose reads the connection's next request again SendReply(connection, finished.reply, nullptr); if (const auto timing = g_DeferredTiming.find(finished.connection); timing != g_DeferredTiming.end()) { - CountRequest(timing->second.route, static_cast(finished.reply.status), timing->second.started, ReplyBytes(finished.reply)); + CountRequest(timing->second.route, static_cast(finished.reply.status), timing->second.started, ReplyBytes(finished.reply), timing->second.fromServer, + timing->second.address, timing->second.requestBytes); g_DeferredTiming.erase(timing); } RemoveSentFile(finished.reply); diff --git a/dWorldServer/WorldServer.cpp b/dWorldServer/WorldServer.cpp index 8de3e57cd..592b180b8 100644 --- a/dWorldServer/WorldServer.cpp +++ b/dWorldServer/WorldServer.cpp @@ -327,6 +327,17 @@ int main(int argc, char** argv) { zoneID); WorldMigration::SetCleanupHandler(CleanupDisconnectedUser); DashboardActions::SetLogoutHandler(CleanupDisconnectedUser); + // The network page's per-connection list says which player is on a connection (main thread, with the report) + Game::server->SetConnectionIdentity([](const SystemAddress& sysAddr, TrafficStats::Connection& connection) { + auto* const user = UserManager::Instance()->GetUser(sysAddr); + if (!user) return; + connection.accountId = user->GetAccountID(); + connection.account = user->GetUsername(); + if (auto* const character = user->GetLastUsedChar()) { + connection.characterId = static_cast(character->GetObjectID()); + connection.character = character->GetName(); + } + }); //Connect to the chat server: uint32_t chatPort = GeneralUtils::TryParse(Game::config->GetValue("chat_server_port")).value_or(1501); @@ -543,6 +554,7 @@ int main(int argc, char** argv) { //Handle our chat packets: packet = Game::chatServer->Receive(); while (packet) { + ChatServerLink::CountReceived(packet->data, packet->length); HandlePacketChat(packet); Game::chatServer->DeallocatePacket(packet); packet = Game::chatServer->Receive(); @@ -1432,10 +1444,7 @@ namespace { if (lastChar) objectID = lastChar->GetObjectID(); } - const auto routed = ToChat(objectID); - RakNet::BitStream bitStream; - routed.WritePacket(bitStream); - Game::chatServer->Send(&bitStream, SYSTEM_PRIORITY, RELIABLE_ORDERED, 0, Game::chatSysAddr, false); + ChatServerLink::Send(ToChat(objectID), SYSTEM_PRIORITY, RELIABLE_ORDERED); } }; diff --git a/tests/dCommonTests/TrafficStatsTests.cpp b/tests/dCommonTests/TrafficStatsTests.cpp index 420cddcb2..5be94a8cf 100644 --- a/tests/dCommonTests/TrafficStatsTests.cpp +++ b/tests/dCommonTests/TrafficStatsTests.cpp @@ -226,3 +226,44 @@ TEST(TrafficStatsTest, CountingIsCheap) { std::printf("[ ] KeyOf + Recorder::Packet: %lld ns per packet\n", static_cast(ns)); EXPECT_LT(ns, 2000); // generous for slow CI machines and sanitizers } + +TEST(TrafficStatsTest, RecorderSplitsPacketsByPeer) { + Recorder r; + const MessageKey in{ false, 4, 5, 0 }; + const MessageKey out{ true, 5, 12, 0 }; + r.Packet(2000, in, 100); // clients by default + r.Packet(2000, out, 40, 3, Peer::CLIENTS); + r.Packet(2000, out, 20, 1, Peer::MASTER); + r.Packet(2000, in, 60, 1, Peer::MASTER); + r.Packet(2000, in, 8, 1, Peer::SERVERS); + const auto report = r.Take(2001); + EXPECT_TRUE(report.peerSplit); + ASSERT_EQ(report.seconds.size(), 1u); + const auto& s = report.seconds[0]; + EXPECT_EQ(s.peers[0], (PeerCounts{ 1, 3, 100, 120 })); + EXPECT_EQ(s.peers[1], (PeerCounts{ 1, 1, 60, 20 })); + EXPECT_EQ(s.peers[2], (PeerCounts{ 1, 0, 8, 0 })); + // The split adds up to the totals + uint64_t packetsIn = 0, bytesOut = 0; + for (const auto& p : s.peers) { packetsIn += p.packetsIn; bytesOut += p.bytesOut; } + EXPECT_EQ(packetsIn, s.packetsIn); + EXPECT_EQ(bytesOut, s.bytesOut); + + Second merged = s; + merged.Merge(s); + EXPECT_EQ(merged.peers[1], (PeerCounts{ 2, 2, 120, 40 })); +} + +TEST(TrafficStatsTest, RecorderSplitsHttpByWhoAsked) { + Recorder r; + r.Http(3000, "GET /api/a", 200, 100, 1000); + r.Http(3000, "GET /api/a", 200, 100, 500, true); + r.HttpOut(3000, 700); + const auto report = r.Take(3001); + ASSERT_EQ(report.seconds.size(), 1u); + EXPECT_EQ(report.seconds[0].httpRequests, 2u); + EXPECT_EQ(report.seconds[0].httpFromServers, 1u); + EXPECT_EQ(report.seconds[0].httpFromServersBytesOut, 500u); + EXPECT_EQ(report.seconds[0].httpOutRequests, 1u); + EXPECT_EQ(report.seconds[0].httpOutBytesIn, 700u); +}