From 6f334c2e3df7f33c2ac2631237083b96644302af Mon Sep 17 00:00:00 2001 From: Aaron Kimbrell Date: Sun, 27 Sep 2026 08:44:03 -0500 Subject: [PATCH] feat(web): record HTTP requests by route with latency for the traffic report Requests are counted by route pattern and status class with a latency histogram (deferred requests until their answer goes out) and bytes sent. The web server reports pending deferred requests and WebSocket clients, the UGC server its worker threads. Co-Authored-By: Claude Opus 5.5 --- dUgcServer/UgcProcessor.h | 5 +++++ dUgcServer/UgcServer.cpp | 8 +++++++ dWeb/Web.cpp | 46 +++++++++++++++++++++++++++++++++++++++ 3 files changed, 59 insertions(+) diff --git a/dUgcServer/UgcProcessor.h b/dUgcServer/UgcProcessor.h index b4d9bf727..e3fa4ba15 100644 --- a/dUgcServer/UgcProcessor.h +++ b/dUgcServer/UgcProcessor.h @@ -58,6 +58,11 @@ public: // Main thread: what the server is doing, for /status nlohmann::json Status() const; + // Any thread: jobs waiting for a worker, and workers busy (traffic diagnostics) + size_t Queued() const { std::lock_guard lock(m_Mutex); return m_Jobs.size(); } + size_t Busy() const { std::lock_guard lock(m_Mutex); return m_Active; } + size_t Threads() const { return m_Config.threads; } + private: struct Job { Kind kind{}; diff --git a/dUgcServer/UgcServer.cpp b/dUgcServer/UgcServer.cpp index 611cc4f92..222fa88d7 100644 --- a/dUgcServer/UgcServer.cpp +++ b/dUgcServer/UgcServer.cpp @@ -23,6 +23,7 @@ #include "Web.h" #include "dConfig.h" #include "dServer.h" +#include "TrafficStats.h" #include "eHTTPMethod.h" #include "json.hpp" @@ -317,6 +318,10 @@ int main(int argc, char** argv) { processorConfig.threads = threads > 0 ? threads : std::max(std::thread::hardware_concurrency() / 2, 1); UgcProcessor processor(processorConfig, storage, library, ReadSettings()); g_Processor = &processor; + // Sent with the traffic reports to the dashboard (Diagnostics) + TrafficStats::Local().SetGauge("workers_busy", [&processor] { return static_cast(processor.Busy()); }); + TrafficStats::Local().SetGauge("workers_queued", [&processor] { return static_cast(processor.Queued()); }); + TrafficStats::Local().SetGauge("workers_threads", [&processor] { return static_cast(processor.Threads()); }); const auto listenIp = Game::config->GetValue("listen_ip").empty() ? std::string("0.0.0.0") : Game::config->GetValue("listen_ip"); const auto port = Setting("port", 2008); @@ -346,6 +351,9 @@ int main(int argc, char** argv) { LOG("Stopping the UGC server"); processor.Stop(); + TrafficStats::Local().SetGauge("workers_busy", nullptr); + TrafficStats::Local().SetGauge("workers_queued", nullptr); + TrafficStats::Local().SetGauge("workers_threads", nullptr); Game::web.Shutdown(); g_Processor = nullptr; Database::Destroy("UgcServer"); diff --git a/dWeb/Web.cpp b/dWeb/Web.cpp index 9c0e19ca0..fbd81305e 100644 --- a/dWeb/Web.cpp +++ b/dWeb/Web.cpp @@ -13,6 +13,9 @@ #include #include #include +#include +#include +#include "TrafficStats.h" namespace Game { Web web; @@ -231,10 +234,39 @@ static void SendReply(mg_connection* connection, const HTTPReply& reply, const m connection->is_resp = 0; } +namespace { + using TrafficClock = std::chrono::steady_clock; + + // Deferred requests' route and start, so their latency counts once the answer goes out + struct DeferredTiming { + std::string route; + TrafficClock::time_point started; + }; + std::unordered_map g_DeferredTiming; + + // The bytes a reply puts on the wire besides its headers (a served file: its size) + uint64_t ReplyBytes(const HTTPReply& reply) { + if (!reply.file.empty() && reply.status == eHTTPStatusCode::OK) { + std::error_code ec; + const auto size = std::filesystem::file_size(reply.file, ec); + return ec ? 0 : size; + } + return reply.message.size(); + } + + void CountRequest(const std::string& route, uint16_t status, TrafficClock::time_point started, uint64_t bytes) { + 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); + } +} + void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_msg) { if (g_HTTPRoutes.empty()) return; HTTPReply reply; + 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)"; if (!http_msg) { reply.status = eHTTPStatusCode::BAD_REQUEST; @@ -321,6 +353,7 @@ void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_ms } } + CountRequest("GET /ws", level ? 101 : 401, started, 0); if (level) { mg_ws_upgrade(connection, const_cast(http_msg), NULL); g_AuthenticatedWSConnections[connection] = { level->level, level->accountId, connectToken, apiToken, @@ -403,6 +436,7 @@ void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_ms if (routeItr != g_HTTPRoutes.end()) { const auto& route = routeItr->second; + trafficRoute = method_string + " " + routeItr->first.second; // Create HTTP context from request HTTPContext context; @@ -449,12 +483,16 @@ void HandleHTTPMessage(mg_connection* connection, const mg_http_message* http_ms // Answered later (Web::Defer): the connection keeps is_resp set, so mongoose reads no further request on it const auto* cc = http_msg ? mg_http_get_header(const_cast(http_msg), "Connection") : nullptr; 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 }; 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)); } @@ -708,6 +746,10 @@ bool Web::Startup(const std::string& listen_ip, const uint32_t listen_port) { .handle = HandleWSGetSubscriptions }); + // Sent with the traffic reports (dServer) + TrafficStats::Local().SetGauge("http_deferred_pending", [] { return static_cast(g_Deferred.Pending()); }); + TrafficStats::Local().SetGauge("websocket_clients", [] { return static_cast(g_AuthenticatedWSConnections.size()); }); + return true; } @@ -724,6 +766,10 @@ void Web::SendDeferredReplies() { if (!connection || connection->is_closing) continue; // 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)); + g_DeferredTiming.erase(timing); + } if (finished.close) connection->is_draining = 1; } }