mirror of
https://github.com/DarkflameUniverse/DarkflameServer.git
synced 2026-10-02 02:43:44 +00:00
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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{};
|
||||
|
||||
@@ -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<size_t>(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<double>(processor.Busy()); });
|
||||
TrafficStats::Local().SetGauge("workers_queued", [&processor] { return static_cast<double>(processor.Queued()); });
|
||||
TrafficStats::Local().SetGauge("workers_threads", [&processor] { return static_cast<double>(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<uint32_t>("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");
|
||||
|
||||
46
dWeb/Web.cpp
46
dWeb/Web.cpp
@@ -13,6 +13,9 @@
|
||||
#include <vector>
|
||||
#include <cctype>
|
||||
#include <chrono>
|
||||
#include <filesystem>
|
||||
#include <unordered_map>
|
||||
#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<unsigned long, DeferredTiming> 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<std::chrono::microseconds>(TrafficClock::now() - started).count();
|
||||
TrafficStats::Local().Http(TrafficStats::Now(), route, status, static_cast<uint64_t>(std::max<int64_t>(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<mg_http_message*>(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<mg_http_message*>(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<uint16_t>(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<double>(g_Deferred.Pending()); });
|
||||
TrafficStats::Local().SetGauge("websocket_clients", [] { return static_cast<double>(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<uint16_t>(finished.reply.status), timing->second.started, ReplyBytes(finished.reply));
|
||||
g_DeferredTiming.erase(timing);
|
||||
}
|
||||
if (finished.close) connection->is_draining = 1;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user