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); +}