From ad1d3b3ac90a50ed6196448c4549875ea0fcd443 Mon Sep 17 00:00:00 2001 From: zhenyan121 <3367366583@qq.com> Date: Wed, 1 Jul 2026 14:53:23 +0800 Subject: [PATCH] feat(networking): add priority and sequence ordering to packet send queues --- include/Cubed/gameplay/network_client.hpp | 33 +++++++++++++++++++---- include/Cubed/gameplay/session.hpp | 32 +++++++++++++++++++--- src/gameplay/client_world.cpp | 6 ++--- src/gameplay/network_client.cpp | 14 +++++----- src/gameplay/server_world.cpp | 16 +++++------ src/gameplay/session.cpp | 14 +++++----- 6 files changed, 83 insertions(+), 32 deletions(-) diff --git a/include/Cubed/gameplay/network_client.hpp b/include/Cubed/gameplay/network_client.hpp index afc8073..6f18dc9 100644 --- a/include/Cubed/gameplay/network_client.hpp +++ b/include/Cubed/gameplay/network_client.hpp @@ -3,6 +3,7 @@ #include "Cubed/gameplay/packet.hpp" #include +#include #include #include namespace Cubed { @@ -14,28 +15,50 @@ public: ~NetworkClient(); void close(); void stop(); - void send(Packet packet); + void send(Packet packet, int priority = 10); void start(std::string ip, int port = 25530); bool is_connected() const; bool is_connect_error() const; private: + struct Task { + int priority = 10; + std::uint64_t sequence = 0; + Packet packet; + Task(int p, std::uint64_t seq, Packet pac) + : priority(p), sequence(seq), packet(std::move(pac)) {} + }; + + struct TaskCompare { + bool operator()(const Task& a, const Task& b) const { + + if (a.priority != b.priority) { + return a.priority > b.priority; + } + + return a.sequence > b.sequence; + } + }; + asio::io_context m_io; std::thread m_net_thread; static constexpr uint32_t MAX_PACKET_SIZE = 4 * 1024 * 1024; tcp::socket m_socket; std::vector m_read_buffer; - std::deque m_write_queue; - asio::strand m_strand; - asio::awaitable connect(std::string ip, int port); - asio::awaitable read_loop(); + std::priority_queue, TaskCompare> m_write_queue; + + asio::strand m_strand; std::atomic m_closed{false}; std::atomic m_connected{false}; std::atomic m_connect_error{false}; // ClientWorld is managed by App ClientWorld& m_world; + std::atomic_uint64_t m_sequence{0}; + + asio::awaitable connect(std::string ip, int port); + asio::awaitable read_loop(); void do_write(); }; diff --git a/include/Cubed/gameplay/session.hpp b/include/Cubed/gameplay/session.hpp index aa4a5da..8e3153b 100644 --- a/include/Cubed/gameplay/session.hpp +++ b/include/Cubed/gameplay/session.hpp @@ -3,8 +3,8 @@ #include "Cubed/gameplay/packet.hpp" #include -#include #include +#include #include namespace Cubed { @@ -17,20 +17,44 @@ public: asio::io_context& io); ~Session(); void start(); - void send(Packet packet); + void send(Packet packet, int priority = 10); + void close(); const std::string& uuid() const; private: + struct Task { + int priority = 10; + std::uint64_t sequence = 0; + Packet packet; + Task(int p, std::uint64_t seq, Packet pac) + : priority(p), sequence(seq), packet(std::move(pac)) {} + }; + + struct TaskCompare { + bool operator()(const Task& a, const Task& b) const { + + if (a.priority != b.priority) { + return a.priority > b.priority; + } + + return a.sequence > b.sequence; + } + }; + static constexpr uint32_t MAX_PACKET_SIZE = 4 * 1024 * 1024; tcp::socket m_socket; std::vector m_read_buffer; - std::deque m_write_queue; + std::priority_queue, TaskCompare> m_write_queue; asio::strand m_strand; std::string m_uuid; ServerWorld& m_server_world; - asio::awaitable read_loop(); std::atomic m_closed{false}; + + std::atomic_uint64_t m_sequence{0}; + + asio::awaitable read_loop(); + void do_write(); }; } // namespace Cubed diff --git a/src/gameplay/client_world.cpp b/src/gameplay/client_world.cpp index 633eb83..987a049 100644 --- a/src/gameplay/client_world.cpp +++ b/src/gameplay/client_world.cpp @@ -270,7 +270,7 @@ void ClientWorld::report_block_change(const glm::ivec3& pos, p->set_x(pos.x); p->set_y(pos.y); p->set_z(pos.z); - m_client->send(make_packet(*req)); + m_client->send(make_packet(*req), 0); } void ClientWorld::receive_block_change(const BlockChangeRsp& rsp) { @@ -336,7 +336,7 @@ void ClientWorld::init(std::string_view player_name, start_thread_pool(); // request login Logger::info("Send Login Request"); - m_client->send(make_packet(req)); + m_client->send(make_packet(req), 0); } void ClientWorld::start_client_thread(std::string_view uuid) { @@ -427,7 +427,7 @@ void ClientWorld::report_player_pos() { v3->set_x(player_pos.x); v3->set_y(player_pos.y); v3->set_z(player_pos.z); - m_client->send(make_packet(*pos)); + m_client->send(make_packet(*pos), 0); } void ClientWorld::update_chunk(const ChunkPosSet& old, const ChunkPosSet& now) { diff --git a/src/gameplay/network_client.cpp b/src/gameplay/network_client.cpp index f47f068..e74b441 100644 --- a/src/gameplay/network_client.cpp +++ b/src/gameplay/network_client.cpp @@ -142,14 +142,15 @@ asio::awaitable NetworkClient::read_loop() { co_return; } -void NetworkClient::send(std::shared_ptr> packet) { +void NetworkClient::send(Packet packet, int priority) { if (m_closed.load()) { return; } - asio::post(m_strand, [self = shared_from_this(), - packet = std::move(packet)]() mutable { + asio::post(m_strand, [self = shared_from_this(), packet = std::move(packet), + priority]() mutable { bool idle = self->m_write_queue.empty(); - self->m_write_queue.emplace_back(std::move(packet)); + self->m_write_queue.emplace(priority, self->m_sequence++, + std::move(packet)); if (idle) { self->do_write(); } @@ -162,15 +163,16 @@ void NetworkClient::do_write() { } auto self = shared_from_this(); + auto packet = std::move(m_write_queue.top().packet); asio::async_write( - m_socket, asio::buffer(*(m_write_queue.front())), + m_socket, asio::buffer(*packet), asio::bind_executor(m_strand, [self](std::error_code ec, size_t) { if (ec) { Logger::warn("Write Ec {}", ec.message()); self->close(); return; } - self->m_write_queue.pop_front(); + self->m_write_queue.pop(); if (!self->m_write_queue.empty()) { self->do_write(); } diff --git a/src/gameplay/server_world.cpp b/src/gameplay/server_world.cpp index c587d64..c3e10cb 100644 --- a/src/gameplay/server_world.cpp +++ b/src/gameplay/server_world.cpp @@ -84,7 +84,7 @@ void ServerWorld::send_time() { rsp->set_game_tick(m_game_ticks); for (auto& [uuid, player] : m_players) { - player.get_session()->send(make_packet(*rsp)); + player.get_session()->send(make_packet(*rsp), 3); } } @@ -573,7 +573,7 @@ void ServerWorld::sync_player_pos(const std::string& uuid, float x, float y, pos->set_x(x); pos->set_y(y); pos->set_z(z); - session->send(make_packet(*rsp)); + session->send(make_packet(*rsp), 0); } } @@ -597,7 +597,7 @@ void ServerWorld::handle_player_login(const std::string& name, if (!sucess) { auto* rsp = Arena::Create(&arena); rsp->set_success(false); - session->send(make_packet(*rsp)); + session->send(make_packet(*rsp), 0); return; } @@ -621,7 +621,7 @@ void ServerWorld::handle_player_login(const std::string& name, auto* rsp = Arena::Create(&arena); rsp->set_success(true); rsp->set_uuid(uuid); - session->send(make_packet(*rsp)); + session->send(make_packet(*rsp), 0); } void ServerWorld::handle_player_exit(const std::string& uuid) { @@ -649,7 +649,7 @@ void ServerWorld::handle_player_exit(const std::string& uuid) { auto* rsp = Arena::Create(&arena); rsp->set_uuid(uuid); rsp->set_server_stop(false); - exit_session->send(make_packet(*rsp)); + exit_session->send(make_packet(*rsp), 0); std::vector> sessions; { @@ -661,7 +661,7 @@ void ServerWorld::handle_player_exit(const std::string& uuid) { for (auto& s : sessions) { if (s) { - s->send(make_packet(*rsp)); + s->send(make_packet(*rsp), 0); } } } @@ -723,7 +723,7 @@ void ServerWorld::handle_block_change(const BlockChangeReq& req) { for (auto& x : sessions) { if (x) { - x->send(make_packet(*rsp)); + x->send(make_packet(*rsp), 1); } } } @@ -783,7 +783,7 @@ void ServerWorld::send_server_stop() { rsp->set_server_stop(true); std::shared_lock lock(m_player_mutex); for (auto& [uuid, player] : m_players) { - player.get_session()->send(make_packet(*rsp)); + player.get_session()->send(make_packet(*rsp), 0); } Logger::info("Send Server Mesaage Success"); } diff --git a/src/gameplay/session.cpp b/src/gameplay/session.cpp index 29dc146..496aca1 100644 --- a/src/gameplay/session.cpp +++ b/src/gameplay/session.cpp @@ -21,11 +21,12 @@ void Session::start() { asio::detached); } -void Session::send(std::shared_ptr> packet) { - asio::post(m_strand, [self = shared_from_this(), - packet = std::move(packet)]() mutable { +void Session::send(std::shared_ptr> packet, int priority) { + asio::post(m_strand, [self = shared_from_this(), packet = std::move(packet), + priority]() mutable { bool idle = self->m_write_queue.empty(); - self->m_write_queue.emplace_back(std::move(packet)); + self->m_write_queue.emplace(priority, self->m_sequence++, + std::move(packet)); if (idle) { self->do_write(); } @@ -118,15 +119,16 @@ asio::awaitable Session::read_loop() { void Session::do_write() { auto self = shared_from_this(); + auto packet = std::move(m_write_queue.top().packet); asio::async_write( - m_socket, asio::buffer(*(m_write_queue.front())), + m_socket, asio::buffer(*packet), asio::bind_executor(m_strand, [self](std::error_code ec, size_t) { if (ec) { Logger::warn("Write Ec {}", ec.message()); self->close(); return; } - self->m_write_queue.pop_front(); + self->m_write_queue.pop(); if (!self->m_write_queue.empty()) { self->do_write(); }