mirror of
https://github.com/DarkflameUniverse/DarkflameServer.git
synced 2026-10-02 10:53:44 +00:00
Every MASTER service packet is now an LUBitStream struct (docs/PacketArchitecture.md, PR 14) and the master server's switch is a dispatch map (PacketDispatcher), as are the master handlers of the world, chat and dashboard servers. - dNet/MasterPackets.h: RequestZoneTransfer, RequestZoneTransferResponse, ServerInfo, RequestSessionKey, SetSessionKey, SessionKeyResponse, NewSessionAlert, PlayerAdded / PlayerRemoved, CreatePrivateZone, RequestPrivateZone (passwords still cut to 50 characters when read), WorldReady, WorldReadyInfo (WORLD_READY to the dashboard), PrepZone, Shutdown, ShutdownResponse, WorldShutDown (SHUTDOWN_RESPONSE to the dashboard), ShutdownUniverse, AffirmTransferRequest/Response, RequestServerList, ServerListResponse, DashboardShutdown, ConfigReload, InstanceShutdown. The Send* functions are gone; MasterPackets::SendToMaster(msg) and SendTo(sysAddr, msg) send a struct. - The dashboard and instance migration structs (PlayerAction, DataChanged, Dashboard messages, MessageCapture, InstanceMigration) are LUBitStreams of the MASTER service now and moved to dNet/master/, included by MasterPackets.h. Their payloads are unchanged; master forwards them by re-serializing the struct instead of copying raw bytes. - InstanceManager, ZoneInstanceManager, MigrationCoordinator, dServer (server info, zone transfer response), auth (SET_SESSION_KEY), the world (session keys, player added and removed, world ready, shutdown response, affirmations, prep zone, shutdown universe) and the dashboard (server list, instance shutdown, config reload, announcements, player actions, message capture) send and read structs. - The login stamps are a `stamps` field of RequestZoneTransfer and RequestZoneTransferResponse (read leniently as before: a message without them reads as empty); master adds its stamps in the REQUEST_ZONE_TRANSFER handler and when it answers, as it did. - InstanceManager::GetInstanceBySysAddr takes a const address. Verified: tests/dGameTests/dNetTests/Legacy/MasterPacketsLegacy.h is a verbatim copy of the old writers and readers; MasterPacketsTests requires identical bytes for a grid of inputs, checks the old readers read what the structs write, round trips and truncation, hand written golden packets, that zone transfers without stamps still read, that the dashboard/migration structs write what "header + Serialize" wrote, and that the dispatcher drops truncated packets. No wire bytes changed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
214 lines
8.5 KiB
C++
214 lines
8.5 KiB
C++
#include "MessageInspector.h"
|
|
#include "MasterPackets.h"
|
|
|
|
#include <chrono>
|
|
#include <cstring>
|
|
#include <optional>
|
|
#include <vector>
|
|
|
|
#include "BitStream.h"
|
|
#include "BitStreamUtils.h"
|
|
#include "Entity.h"
|
|
#include "Game.h"
|
|
#include "GameMessageDecoder.h"
|
|
#include "Logger.h"
|
|
#include "master/MessageCapture.h"
|
|
#include "MessageType/Client.h"
|
|
#include "MessageType/Master.h"
|
|
#include "PlayerManager.h"
|
|
#include "ServiceType.h"
|
|
#include "dServer.h"
|
|
#include "dZoneManager.h"
|
|
|
|
namespace {
|
|
using Clock = std::chrono::steady_clock;
|
|
|
|
constexpr size_t MAX_CAPTURES = 8;
|
|
// Per capture: at most this many messages a second, and this many waiting to be sent
|
|
constexpr uint32_t MAX_PER_SECOND = 200;
|
|
constexpr size_t MAX_PENDING = 1000;
|
|
constexpr auto SEND_INTERVAL = std::chrono::milliseconds(250);
|
|
// Sent even when nothing was captured, so the dashboard knows the capture is still running
|
|
constexpr auto HEARTBEAT_INTERVAL = std::chrono::seconds(5);
|
|
constexpr size_t BATCH_ENTRIES = 100;
|
|
constexpr size_t BATCH_BYTES = 48 * 1024;
|
|
constexpr size_t BATCHES_PER_SEND = 2;
|
|
|
|
// Game messages to a client start with the packet header, the object ID and the message ID
|
|
constexpr size_t GAME_MESSAGE_HEADER_BYTES = 8 + sizeof(LWOOBJID) + sizeof(MessageType::Game);
|
|
|
|
struct Capture {
|
|
MessageCaptureControl control;
|
|
SystemAddress sysAddr;
|
|
Clock::time_point until;
|
|
MessageCaptureQueue queue{ MAX_PENDING, MAX_PER_SECOND };
|
|
uint32_t nextSequence{ 1 };
|
|
Clock::time_point lastSent{};
|
|
};
|
|
|
|
std::vector<Capture> g_Captures;
|
|
|
|
int64_t NowMs() {
|
|
return std::chrono::duration_cast<std::chrono::milliseconds>(std::chrono::system_clock::now().time_since_epoch()).count();
|
|
}
|
|
|
|
void Send(const Capture& capture, eMessageCaptureStatus status, std::vector<MessageCaptureEntry> entries = {}, uint32_t dropped = 0, eMessageCaptureEnd reason = {}) {
|
|
if (!Game::server) return;
|
|
MessageCaptureData data;
|
|
data.captureId = capture.control.captureId;
|
|
data.status = status;
|
|
data.characterId = capture.control.characterId;
|
|
data.zoneId = Game::zoneManager ? Game::zoneManager->GetZoneID().GetMapID() : 0;
|
|
data.instanceId = Game::server->GetInstanceID();
|
|
data.cloneId = Game::zoneManager ? Game::zoneManager->GetZoneID().GetCloneID() : 0;
|
|
data.reason = reason;
|
|
data.dropped = dropped;
|
|
data.entries = std::move(entries);
|
|
|
|
MasterPackets::SendToMaster(data);
|
|
}
|
|
|
|
// Send what is waiting: a couple of batches per call (the rest goes next time), or everything when `all`
|
|
void Flush(Capture& capture, Clock::time_point now, bool all) {
|
|
if (!all && now - capture.lastSent < SEND_INTERVAL) return;
|
|
bool sent = false;
|
|
for (size_t i = 0; all || i < BATCHES_PER_SEND; i++) {
|
|
auto batch = capture.queue.Take(BATCH_ENTRIES, BATCH_BYTES);
|
|
const auto dropped = capture.queue.TakeDropped();
|
|
if (batch.empty() && dropped == 0) break;
|
|
Send(capture, eMessageCaptureStatus::ENTRIES, std::move(batch), dropped);
|
|
sent = true;
|
|
}
|
|
if (!sent && !all && now - capture.lastSent >= HEARTBEAT_INTERVAL) {
|
|
Send(capture, eMessageCaptureStatus::ENTRIES);
|
|
sent = true;
|
|
}
|
|
if (sent) capture.lastSent = now;
|
|
}
|
|
|
|
void SetHook();
|
|
|
|
void End(size_t index, eMessageCaptureEnd reason) {
|
|
auto& capture = g_Captures[index];
|
|
Flush(capture, Clock::now(), true);
|
|
Send(capture, eMessageCaptureStatus::ENDED, {}, 0, reason);
|
|
LOG("Message capture %u of character %llu ended (%i)", capture.control.captureId, capture.control.characterId, static_cast<int>(reason));
|
|
g_Captures.erase(g_Captures.begin() + index);
|
|
SetHook();
|
|
}
|
|
|
|
/**
|
|
* Add a message to every capture that wants it. Decoding happens once, and only for messages a capture keeps.
|
|
* `address` is who the message goes to (or, for a broadcast, who it skips).
|
|
*/
|
|
void Record(eMessageDirection direction, const SystemAddress& address, bool broadcast, LWOOBJID objectId, uint16_t messageId,
|
|
const unsigned char* data, uint32_t bits) {
|
|
const auto now = Clock::now();
|
|
const auto timeMs = NowMs();
|
|
std::optional<std::string> decoded;
|
|
for (auto& capture : g_Captures) {
|
|
// A broadcast reaches everyone except the address given (when one is given)
|
|
const bool reaches = broadcast ? capture.sysAddr != address : capture.sysAddr == address;
|
|
if (!reaches || !capture.control.Wants(direction, messageId)) continue;
|
|
|
|
const auto bytes = (bits + 7) / 8;
|
|
if (!decoded) {
|
|
decoded.emplace();
|
|
const bool toServer = direction == eMessageDirection::TO_SERVER;
|
|
if (GameMessageDecoder::CanDecode(static_cast<MessageType::Game>(messageId), toServer)) {
|
|
RakNet::BitStream payload(const_cast<unsigned char*>(data), bytes, true);
|
|
const auto fields = GameMessageDecoder::Decode(static_cast<MessageType::Game>(messageId), toServer, payload);
|
|
if (fields) {
|
|
*decoded = fields->dump();
|
|
if (decoded->size() > MessageCapture::MAX_DECODED) decoded->clear();
|
|
}
|
|
}
|
|
}
|
|
|
|
MessageCaptureEntry entry;
|
|
entry.sequence = capture.nextSequence++;
|
|
entry.timeMs = timeMs;
|
|
entry.direction = direction;
|
|
entry.messageId = messageId;
|
|
entry.objectId = objectId;
|
|
entry.bits = bits;
|
|
entry.payload.assign(reinterpret_cast<const char*>(data), std::min<size_t>(bytes, MessageCapture::MAX_PAYLOAD));
|
|
entry.decoded = *decoded;
|
|
capture.queue.Push(std::move(entry), now);
|
|
}
|
|
}
|
|
|
|
// Packets the server sends: only game messages are kept
|
|
void OnSend(const RakNet::BitStream& bitStream, const SystemAddress& sysAddr, bool broadcast) {
|
|
static const auto header = [] {
|
|
RakNet::BitStream expected;
|
|
BitStreamUtils::WriteHeader(expected, ServiceType::CLIENT, MessageType::Client::GAME_MSG);
|
|
return std::string(reinterpret_cast<const char*>(expected.GetData()), expected.GetNumberOfBytesUsed());
|
|
}();
|
|
const auto totalBits = bitStream.GetNumberOfBitsUsed();
|
|
if (totalBits < GAME_MESSAGE_HEADER_BYTES * 8) return;
|
|
const auto* data = bitStream.GetData();
|
|
if (std::memcmp(data, header.data(), header.size()) != 0) return;
|
|
|
|
LWOOBJID objectId{};
|
|
uint16_t messageId{};
|
|
std::memcpy(&objectId, data + 8, sizeof(objectId));
|
|
std::memcpy(&messageId, data + 8 + sizeof(objectId), sizeof(messageId));
|
|
Record(eMessageDirection::TO_CLIENT, sysAddr, broadcast, objectId, messageId, data + GAME_MESSAGE_HEADER_BYTES,
|
|
static_cast<uint32_t>(totalBits - GAME_MESSAGE_HEADER_BYTES * 8));
|
|
}
|
|
|
|
// Watch sent packets only while something is captured
|
|
void SetHook() {
|
|
MessageInspector::g_Capturing = !g_Captures.empty();
|
|
if (!Game::server) return;
|
|
Game::server->SetSendObserver(MessageInspector::g_Capturing ? dServer::SendObserver(OnSend) : nullptr);
|
|
}
|
|
}
|
|
|
|
namespace MessageInspector {
|
|
bool g_Capturing = false;
|
|
|
|
void Control(const MessageCaptureControl& control) {
|
|
const auto existing = std::ranges::find_if(g_Captures, [&](const Capture& c) { return c.control.captureId == control.captureId; });
|
|
if (control.action == eMessageCaptureControl::STOP) {
|
|
if (existing != g_Captures.end()) End(existing - g_Captures.begin(), eMessageCaptureEnd::STOPPED);
|
|
return;
|
|
}
|
|
|
|
// Only the world the character is in captures it
|
|
auto* player = PlayerManager::GetPlayer(control.characterId);
|
|
if (!player || control.seconds == 0) return;
|
|
if (existing == g_Captures.end() && g_Captures.size() >= MAX_CAPTURES) {
|
|
LOG("Not starting message capture %u: %zu captures are already running here", control.captureId, g_Captures.size());
|
|
return;
|
|
}
|
|
|
|
Capture& capture = existing != g_Captures.end() ? *existing : g_Captures.emplace_back();
|
|
capture.control = control;
|
|
capture.sysAddr = player->GetSystemAddress();
|
|
capture.until = Clock::now() + std::chrono::seconds(std::min(control.seconds, MessageCapture::MAX_SECONDS));
|
|
capture.lastSent = Clock::now();
|
|
LOG("Capturing game messages of character %llu for %u seconds (capture %u)", control.characterId, control.seconds, control.captureId);
|
|
Send(capture, eMessageCaptureStatus::STARTED);
|
|
SetHook();
|
|
}
|
|
|
|
void RecordReceived(const SystemAddress& sysAddr, LWOOBJID objectId, MessageType::Game messageId, const RakNet::BitStream& payload) {
|
|
Record(eMessageDirection::TO_SERVER, sysAddr, false, objectId, static_cast<uint16_t>(messageId), payload.GetData(),
|
|
static_cast<uint32_t>(payload.GetNumberOfBitsUsed()));
|
|
}
|
|
|
|
void Update() {
|
|
if (!g_Capturing) return;
|
|
const auto now = Clock::now();
|
|
for (size_t i = g_Captures.size(); i-- > 0;) {
|
|
auto& capture = g_Captures[i];
|
|
const auto* player = PlayerManager::GetPlayer(capture.control.characterId);
|
|
if (!player || player->GetSystemAddress() != capture.sysAddr) End(i, eMessageCaptureEnd::PLAYER_LEFT);
|
|
else if (now >= capture.until) End(i, eMessageCaptureEnd::TIME_LIMIT);
|
|
else Flush(capture, now, false);
|
|
}
|
|
}
|
|
}
|