|
|
|
|
@@ -1,11 +1,14 @@
|
|
|
|
|
#include "CaptureReplay.h"
|
|
|
|
|
#include "MasterPackets.h"
|
|
|
|
|
|
|
|
|
|
#include <atomic>
|
|
|
|
|
#include <chrono>
|
|
|
|
|
#include <ctime>
|
|
|
|
|
#include <deque>
|
|
|
|
|
#include <fstream>
|
|
|
|
|
#include <map>
|
|
|
|
|
#include <memory>
|
|
|
|
|
#include <mutex>
|
|
|
|
|
#include <regex>
|
|
|
|
|
#include <set>
|
|
|
|
|
|
|
|
|
|
@@ -20,6 +23,8 @@
|
|
|
|
|
#include "Logger.h"
|
|
|
|
|
#include "master/MessageCapture.h"
|
|
|
|
|
#include "PacketDecoder.h"
|
|
|
|
|
#include "ReplicaDecoder.h"
|
|
|
|
|
#include "Workers.h"
|
|
|
|
|
#include "Permissions.h"
|
|
|
|
|
#include "RouteUtils.h"
|
|
|
|
|
#include "Web.h"
|
|
|
|
|
@@ -76,14 +81,39 @@ namespace {
|
|
|
|
|
std::map<uint32_t, Capture> g_Captures;
|
|
|
|
|
Clock::time_point g_NextPush{};
|
|
|
|
|
|
|
|
|
|
// LOT -> components, read from the CDClient at startup (PreloadDecoding) and only read after: the replica pass runs
|
|
|
|
|
// on worker threads, which never query the CDClient
|
|
|
|
|
ReplicaDecoder::ComponentTable g_Components;
|
|
|
|
|
|
|
|
|
|
// The last bundle read for the viewer (a running capture's file grows: read again when its size changed)
|
|
|
|
|
struct Cached {
|
|
|
|
|
uint64_t id{};
|
|
|
|
|
uintmax_t size{};
|
|
|
|
|
CaptureBundle::Bundle bundle;
|
|
|
|
|
int64_t startUs{};
|
|
|
|
|
|
|
|
|
|
// Replica packets' fields, read once per capture in timeline order (a serialization needs the construction
|
|
|
|
|
// before it), on a worker thread
|
|
|
|
|
std::once_flag decodeOnce;
|
|
|
|
|
std::atomic<bool> decoded{};
|
|
|
|
|
std::vector<std::optional<nlohmann::json>> replica;
|
|
|
|
|
|
|
|
|
|
void Decode() {
|
|
|
|
|
std::call_once(decodeOnce, [this] {
|
|
|
|
|
ReplicaDecoder::Session session(g_Components);
|
|
|
|
|
replica.resize(bundle.records.size());
|
|
|
|
|
for (size_t i = 0; i < bundle.records.size(); i++) {
|
|
|
|
|
const auto& record = bundle.records[i];
|
|
|
|
|
if (record.header.flags & PacketRecordFlags::GAP || CaptureTools::FromClient(record.header)) continue;
|
|
|
|
|
replica[i] = session.Decode(record.bytes, CaptureTools::ReplicaConnection(record.header));
|
|
|
|
|
}
|
|
|
|
|
decoded = true;
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const nlohmann::json* ReplicaFields(size_t i) const { return i < replica.size() && replica[i] ? &*replica[i] : nullptr; }
|
|
|
|
|
};
|
|
|
|
|
Cached g_Cache;
|
|
|
|
|
std::shared_ptr<Cached> g_Cache;
|
|
|
|
|
|
|
|
|
|
bool Connected() { return Game::server && Game::server->GetIsConnectedToMaster(); }
|
|
|
|
|
|
|
|
|
|
@@ -235,7 +265,7 @@ namespace {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// The capture's records on one timeline (read again when the file grew)
|
|
|
|
|
const Cached* Load(const Session& session, HTTPReply& reply) {
|
|
|
|
|
std::shared_ptr<Cached> Load(const Session& session, HTTPReply& reply) {
|
|
|
|
|
const auto path = CaptureReplay::FileOf(session.id);
|
|
|
|
|
std::error_code ec;
|
|
|
|
|
const auto size = fs::file_size(path, ec);
|
|
|
|
|
@@ -243,8 +273,9 @@ namespace {
|
|
|
|
|
JsonError(reply, eHTTPStatusCode::NOT_FOUND, "The capture's file is gone (deleted, or kept on another machine)");
|
|
|
|
|
return nullptr;
|
|
|
|
|
}
|
|
|
|
|
if (g_Cache.id == session.id && g_Cache.size == size) return &g_Cache;
|
|
|
|
|
Cached cached;
|
|
|
|
|
if (g_Cache && g_Cache->id == session.id && g_Cache->size == size) return g_Cache;
|
|
|
|
|
auto loaded = std::make_shared<Cached>();
|
|
|
|
|
auto& cached = *loaded;
|
|
|
|
|
std::string error;
|
|
|
|
|
if (!CaptureBundle::Load(path, cached.bundle, error)) {
|
|
|
|
|
JsonError(reply, eHTTPStatusCode::INTERNAL_SERVER_ERROR, error);
|
|
|
|
|
@@ -254,8 +285,8 @@ namespace {
|
|
|
|
|
cached.id = session.id;
|
|
|
|
|
cached.size = size;
|
|
|
|
|
cached.startUs = cached.bundle.records.empty() ? session.startedAt * 1000000 : std::min(session.startedAt * 1000000, cached.bundle.records.front().header.timeUs);
|
|
|
|
|
g_Cache = std::move(cached);
|
|
|
|
|
return &g_Cache;
|
|
|
|
|
g_Cache = std::move(loaded);
|
|
|
|
|
return g_Cache;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
std::string CharacterName(LWOOBJID id) {
|
|
|
|
|
@@ -442,6 +473,11 @@ namespace CaptureReplay {
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void PreloadDecoding() {
|
|
|
|
|
const auto lots = ReplicaDecoder::LoadComponentTable(g_Components);
|
|
|
|
|
LOG("Read the components of %zu LOTs for decoding replica packets", lots);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void RegisterRoutes() {
|
|
|
|
|
Game::web.RegisterWSSubscription(TOPIC, std::function<uint8_t()>([] { return Permissions::Level(PERMISSION); }));
|
|
|
|
|
|
|
|
|
|
@@ -574,7 +610,7 @@ namespace CaptureReplay {
|
|
|
|
|
[](HTTPReply& reply, const HTTPContext& context) {
|
|
|
|
|
const auto session = FindSession(context, reply);
|
|
|
|
|
if (!session) return;
|
|
|
|
|
const auto* cached = Load(*session, reply);
|
|
|
|
|
const auto cached = Load(*session, reply);
|
|
|
|
|
if (!cached) return;
|
|
|
|
|
const auto offset = GeneralUtils::TryParse<size_t>(QueryValue(context.queryString, "offset")).value_or(0);
|
|
|
|
|
const auto limit = std::clamp<size_t>(GeneralUtils::TryParse<size_t>(QueryValue(context.queryString, "limit")).value_or(PAGE_DEFAULT), 1, PAGE_MAX);
|
|
|
|
|
@@ -582,39 +618,48 @@ namespace CaptureReplay {
|
|
|
|
|
for (auto& ch : q) ch = static_cast<char>(std::toupper(static_cast<unsigned char>(ch)));
|
|
|
|
|
const auto source = QueryValue(context.queryString, "source");
|
|
|
|
|
const bool fields = QueryValue(context.queryString, "fields") != "0";
|
|
|
|
|
json records = json::array();
|
|
|
|
|
size_t matched = 0;
|
|
|
|
|
const auto& all = cached->bundle.records;
|
|
|
|
|
for (size_t i = 0; i < all.size(); i++) {
|
|
|
|
|
if (q.empty() && source.empty()) {
|
|
|
|
|
if (i < offset) continue;
|
|
|
|
|
if (records.size() >= limit) break;
|
|
|
|
|
records.push_back(CaptureTools::RecordJson(all[i], i, cached->startUs, fields));
|
|
|
|
|
continue;
|
|
|
|
|
// Read on the main thread (the session JSON reads live capture state); the packets are decoded on a
|
|
|
|
|
// worker the first time, since replica packets need a pass over the whole capture
|
|
|
|
|
auto capture = StoredSessionJson(*session);
|
|
|
|
|
Workers::Reply(reply, context, cached->decoded, [cached, capture = std::move(capture), offset, limit, q, source, fields](HTTPReply& out) {
|
|
|
|
|
if (fields) cached->Decode();
|
|
|
|
|
json records = json::array();
|
|
|
|
|
size_t matched = 0;
|
|
|
|
|
const auto& all = cached->bundle.records;
|
|
|
|
|
for (size_t i = 0; i < all.size(); i++) {
|
|
|
|
|
if (q.empty() && source.empty()) {
|
|
|
|
|
if (i < offset) continue;
|
|
|
|
|
if (records.size() >= limit) break;
|
|
|
|
|
records.push_back(CaptureTools::RecordJson(all[i], i, cached->startUs, fields, cached->ReplicaFields(i)));
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
auto record = CaptureTools::RecordJson(all[i], i, cached->startUs, false);
|
|
|
|
|
if (!source.empty() && record.value("source", std::string{}) != source) continue;
|
|
|
|
|
if (!q.empty() && record.value("name", std::string{}).find(q) == std::string::npos) continue;
|
|
|
|
|
if (matched++ < offset || records.size() >= limit) continue;
|
|
|
|
|
records.push_back(fields ? CaptureTools::RecordJson(all[i], i, cached->startUs, true, cached->ReplicaFields(i)) : record);
|
|
|
|
|
}
|
|
|
|
|
auto record = CaptureTools::RecordJson(all[i], i, cached->startUs, false);
|
|
|
|
|
if (!source.empty() && record.value("source", std::string{}) != source) continue;
|
|
|
|
|
if (!q.empty() && record.value("name", std::string{}).find(q) == std::string::npos) continue;
|
|
|
|
|
if (matched++ < offset || records.size() >= limit) continue;
|
|
|
|
|
records.push_back(fields ? CaptureTools::RecordJson(all[i], i, cached->startUs, true) : record);
|
|
|
|
|
}
|
|
|
|
|
JsonSuccess(reply, { {"capture", StoredSessionJson(*session)}, {"total", all.size()}, {"matched", q.empty() && source.empty() ? all.size() : matched},
|
|
|
|
|
{"start", cached->startUs / 1000}, {"duration", all.empty() ? 0.0 : static_cast<double>(all.back().header.timeUs - cached->startUs) / 1000.0},
|
|
|
|
|
{"records", records} });
|
|
|
|
|
JsonSuccess(out, { {"capture", capture}, {"total", all.size()}, {"matched", q.empty() && source.empty() ? all.size() : matched},
|
|
|
|
|
{"start", cached->startUs / 1000}, {"duration", all.empty() ? 0.0 : static_cast<double>(all.back().header.timeUs - cached->startUs) / 1000.0},
|
|
|
|
|
{"records", records} });
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
Route(eHTTPMethod::GET, "/api/inspector/sessions/:id/packets/:index", Perm(PERMISSION), "One captured packet with its bytes",
|
|
|
|
|
[](HTTPReply& reply, const HTTPContext& context) {
|
|
|
|
|
const auto session = FindSession(context, reply);
|
|
|
|
|
if (!session) return;
|
|
|
|
|
const auto* cached = Load(*session, reply);
|
|
|
|
|
const auto cached = Load(*session, reply);
|
|
|
|
|
if (!cached) return;
|
|
|
|
|
const auto index = PathId<size_t>(context.path, 5);
|
|
|
|
|
if (!index || *index >= cached->bundle.records.size()) return JsonError(reply, eHTTPStatusCode::NOT_FOUND, "No such packet");
|
|
|
|
|
const auto& record = cached->bundle.records[*index];
|
|
|
|
|
auto out = CaptureTools::RecordJson(record, *index, cached->startUs, true);
|
|
|
|
|
out["hex"] = MessageCapture::ToHex(record.bytes);
|
|
|
|
|
JsonSuccess(reply, { {"record", out} });
|
|
|
|
|
Workers::Reply(reply, context, cached->decoded, [cached, index = *index](HTTPReply& out) {
|
|
|
|
|
cached->Decode();
|
|
|
|
|
const auto& record = cached->bundle.records[index];
|
|
|
|
|
auto json = CaptureTools::RecordJson(record, index, cached->startUs, true, cached->ReplicaFields(index));
|
|
|
|
|
json["hex"] = MessageCapture::ToHex(record.bytes);
|
|
|
|
|
JsonSuccess(out, { {"record", json} });
|
|
|
|
|
});
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
Route(eHTTPMethod::GET, "/api/inspector/sessions/:id/positions", Perm(PERMISSION),
|
|
|
|
|
@@ -622,7 +667,7 @@ namespace CaptureReplay {
|
|
|
|
|
[](HTTPReply& reply, const HTTPContext& context) {
|
|
|
|
|
const auto session = FindSession(context, reply);
|
|
|
|
|
if (!session) return;
|
|
|
|
|
const auto* cached = Load(*session, reply);
|
|
|
|
|
const auto cached = Load(*session, reply);
|
|
|
|
|
if (!cached) return;
|
|
|
|
|
const auto tracks = CaptureTools::Tracks(cached->bundle.records, cached->startUs);
|
|
|
|
|
std::set<uint32_t> zones;
|
|
|
|
|
@@ -656,7 +701,7 @@ namespace CaptureReplay {
|
|
|
|
|
[](HTTPReply& reply, const HTTPContext& context) {
|
|
|
|
|
const auto session = FindSession(context, reply);
|
|
|
|
|
if (!session) return;
|
|
|
|
|
const auto* cached = Load(*session, reply);
|
|
|
|
|
const auto cached = Load(*session, reply);
|
|
|
|
|
if (!cached) return;
|
|
|
|
|
const bool anonymise = QueryValue(context.queryString, "anonymise") == "1";
|
|
|
|
|
auto bundle = cached->bundle;
|
|
|
|
|
|