mirror of
https://github.com/zhenyan121/Cubed.git
synced 2026-08-08 17:57:02 +08:00
refactor(server): enhance thread safety and session management
This commit is contained in:
@@ -25,6 +25,7 @@ private:
|
||||
int m_port = 25530;
|
||||
std::atomic<bool> m_stopped{false};
|
||||
ServerWorld m_world;
|
||||
std::mutex m_session_mutex;
|
||||
asio::awaitable<void> listen();
|
||||
void net_run();
|
||||
};
|
||||
|
||||
@@ -2,16 +2,20 @@
|
||||
#include "Cubed/gameplay/chunk_pos.hpp"
|
||||
|
||||
#include <glm/glm.hpp>
|
||||
#include <memory>
|
||||
#include <string>
|
||||
#include <string_view>
|
||||
namespace Cubed {
|
||||
class ServerWorld;
|
||||
class Session;
|
||||
class ServerPlayer {
|
||||
public:
|
||||
ServerPlayer(std::string_view name, std::string_view uuid,
|
||||
ServerWorld& m_world);
|
||||
ServerWorld& m_world, std::shared_ptr<Session> session);
|
||||
const glm::vec3& get_pos() const;
|
||||
const std::string& get_name() const;
|
||||
const std::string& get_uuid() const;
|
||||
std::shared_ptr<Session> get_session() const;
|
||||
void update_pos(float x, float y, float z);
|
||||
|
||||
private:
|
||||
@@ -20,5 +24,6 @@ private:
|
||||
glm::vec3 m_pos{0.0f};
|
||||
ServerWorld& m_world;
|
||||
ChunkPos m_last_chunk_pos{0, 0};
|
||||
std::shared_ptr<Session> m_session;
|
||||
};
|
||||
} // namespace Cubed
|
||||
|
||||
@@ -24,7 +24,7 @@ public:
|
||||
ServerWorld();
|
||||
~ServerWorld();
|
||||
void player_join(std::string_view name, std::string_view uuid);
|
||||
void player_exit(const std::string& name);
|
||||
void handle_player_exit(const std::string& name);
|
||||
void init_world();
|
||||
void need_gen(std::optional<std::string> uuid);
|
||||
void update();
|
||||
@@ -81,20 +81,18 @@ private:
|
||||
std::future<void> future;
|
||||
};
|
||||
using ChunkHashMap =
|
||||
std::unordered_map<ChunkPos, ServerChunk, ChunkPos::Hash>;
|
||||
using PlayerHashMap =
|
||||
tbb::concurrent_unordered_map<std::string, ServerPlayer>;
|
||||
tbb::concurrent_unordered_map<ChunkPos, ServerChunk, ChunkPos::Hash>;
|
||||
using PlayerHashMap = std::unordered_map<std::string, ServerPlayer>;
|
||||
using PendingChunkHashMap =
|
||||
std::unordered_map<ChunkPos, PendingChunk, ChunkPos::Hash>;
|
||||
using ChunkPosSet = std::unordered_set<ChunkPos, ChunkPos::Hash>;
|
||||
using PlayerSessionMap =
|
||||
tbb::concurrent_hash_map<std::string, std::shared_ptr<Session>>;
|
||||
using session_acc = PlayerSessionMap::accessor;
|
||||
using session_cacc = PlayerSessionMap::const_accessor;
|
||||
using PlayerUUIDMap = tbb::concurrent_hash_map<std::string, std::string>;
|
||||
using uuid_acc = PlayerUUIDMap::accessor;
|
||||
using uuid_cacc = PlayerUUIDMap::const_accessor;
|
||||
// key = uuid
|
||||
PlayerHashMap m_players;
|
||||
ChunkHashMap m_chunks;
|
||||
// Can only be used in the gen thread
|
||||
PendingChunkHashMap new_chunks;
|
||||
PendingChunkHashMap m_new_chunks;
|
||||
std::vector<std::pair<ChunkPos, ServerChunk>> m_new_finished_chunk;
|
||||
|
||||
CaveCarver m_cave_carcer;
|
||||
@@ -119,8 +117,9 @@ private:
|
||||
|
||||
std::shared_mutex m_chunks_mutex;
|
||||
std::shared_mutex m_new_chunk_mutex;
|
||||
std::shared_mutex m_player_mutex;
|
||||
std::mutex m_need_gen_queue_mutex;
|
||||
std::condition_variable m_gen_cv;
|
||||
std::condition_variable_any m_gen_cv;
|
||||
|
||||
std::deque<std::string> m_need_gen_queue;
|
||||
|
||||
@@ -128,9 +127,7 @@ private:
|
||||
|
||||
std::atomic<ChunkLoadStyle> m_chunk_load_style{ChunkLoadStyle::RANDOM};
|
||||
|
||||
// key = uuid
|
||||
PlayerSessionMap m_player_session;
|
||||
tbb::concurrent_hash_map<std::string, std::string> m_uuid_to_name;
|
||||
PlayerUUIDMap m_uuid_to_name;
|
||||
void init_chunks();
|
||||
|
||||
void gen_chunks_internal(std::optional<std::string> uuid);
|
||||
|
||||
@@ -12,14 +12,25 @@ void NetworkServer::stop() {
|
||||
if (m_stopped.exchange(true)) {
|
||||
return;
|
||||
}
|
||||
for (auto& [key, s] : m_session) {
|
||||
if (s) {
|
||||
s->close();
|
||||
}
|
||||
}
|
||||
|
||||
m_io.stop();
|
||||
|
||||
std::vector<std::shared_ptr<Session>> sessions;
|
||||
|
||||
{
|
||||
std::lock_guard lock(m_session_mutex);
|
||||
|
||||
for (auto& [id, s] : m_session) {
|
||||
sessions.push_back(s);
|
||||
}
|
||||
|
||||
m_session.clear();
|
||||
}
|
||||
|
||||
for (auto& s : sessions) {
|
||||
s->close();
|
||||
}
|
||||
|
||||
if (m_net_thread.joinable()) {
|
||||
m_net_thread.join();
|
||||
}
|
||||
@@ -27,22 +38,36 @@ void NetworkServer::stop() {
|
||||
}
|
||||
|
||||
asio::awaitable<void> NetworkServer::listen() {
|
||||
|
||||
try {
|
||||
tcp::acceptor acceptor(m_io, tcp::endpoint(tcp::v4(), m_port));
|
||||
while (true) {
|
||||
tcp::socket socket =
|
||||
co_await acceptor.async_accept(asio::use_awaitable);
|
||||
if (m_stopped) {
|
||||
break;
|
||||
}
|
||||
|
||||
std::shared_ptr<Session> s =
|
||||
std::make_shared<Session>(std::move(socket), m_world);
|
||||
s->start();
|
||||
{
|
||||
std::lock_guard lock(m_session_mutex);
|
||||
m_session.emplace(s->uuid(), s);
|
||||
}
|
||||
s->start();
|
||||
}
|
||||
} catch (const std::exception& e) {
|
||||
if (!m_stopped) {
|
||||
Logger::error("accept error {}", e.what());
|
||||
}
|
||||
} catch (...) {
|
||||
Logger::error("Network Server: Unkown Error");
|
||||
}
|
||||
|
||||
co_return;
|
||||
}
|
||||
|
||||
void NetworkServer::net_run() {
|
||||
if (m_net_thread.joinable()) {
|
||||
return;
|
||||
}
|
||||
m_net_thread = std::thread([this]() {
|
||||
asio::co_spawn(m_io, listen(), asio::detached);
|
||||
m_io.run();
|
||||
|
||||
@@ -3,10 +3,12 @@
|
||||
#include "Cubed/gameplay/server_world.hpp"
|
||||
namespace Cubed {
|
||||
ServerPlayer::ServerPlayer(std::string_view name, std::string_view uuid,
|
||||
ServerWorld& world)
|
||||
: m_name(name), m_uuid(uuid), m_world(world) {}
|
||||
ServerWorld& world, std::shared_ptr<Session> session)
|
||||
: m_name(name), m_uuid(uuid), m_world(world), m_session(session) {}
|
||||
const glm::vec3& ServerPlayer::get_pos() const { return m_pos; }
|
||||
const std::string& ServerPlayer::get_name() const { return m_name; }
|
||||
const std::string& ServerPlayer::get_uuid() const { return m_uuid; }
|
||||
std::shared_ptr<Session> ServerPlayer::get_session() const { return m_session; }
|
||||
void ServerPlayer::update_pos(float x, float y, float z) {
|
||||
m_pos = glm::vec3{x, y, z};
|
||||
ChunkPos chunk_pos = get_chunk_pos(x, z);
|
||||
|
||||
@@ -24,7 +24,8 @@ ServerWorld::~ServerWorld() {
|
||||
}
|
||||
|
||||
void ServerWorld::wait_all_chunk_tasks() {
|
||||
for (auto& [pos, task] : new_chunks) {
|
||||
std::lock_guard lock(m_new_chunk_mutex);
|
||||
for (auto& [pos, task] : m_new_chunks) {
|
||||
task.future.get();
|
||||
}
|
||||
}
|
||||
@@ -71,9 +72,11 @@ void ServerWorld::gen_chunks_internal(std::optional<std::string> uuid) {
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
{
|
||||
std::lock_guard lock(m_new_chunk_mutex);
|
||||
for (auto& pos : need_gen_chunks_pos) {
|
||||
new_chunks.emplace(pos, ServerChunk(*this, pos));
|
||||
m_new_chunks.emplace(pos, ServerChunk(*this, pos));
|
||||
}
|
||||
}
|
||||
|
||||
submit_new_chunks(uuid);
|
||||
@@ -109,7 +112,7 @@ void ServerWorld::sync_and_collect_missing_chunks(
|
||||
std::lock_guard lk(m_chunks_mutex);
|
||||
for (auto it = m_chunks.begin(); it != m_chunks.end();) {
|
||||
if (required_chunks.find(it->first) == required_chunks.end()) {
|
||||
it = m_chunks.erase(it);
|
||||
it = m_chunks.unsafe_erase(it);
|
||||
} else {
|
||||
++it;
|
||||
}
|
||||
@@ -131,7 +134,7 @@ void ServerWorld::submit_new_chunks(const std::optional<std::string>& uuid) {
|
||||
}
|
||||
switch (m_chunk_load_style) {
|
||||
case RANDOM:
|
||||
for (auto& [pos, task] : new_chunks) {
|
||||
for (auto& [pos, task] : m_new_chunks) {
|
||||
if (!task.future.valid()) {
|
||||
task.future =
|
||||
pool_ptr->enqueue([&task]() { task.chunk.gen_chunk(); });
|
||||
@@ -140,7 +143,7 @@ void ServerWorld::submit_new_chunks(const std::optional<std::string>& uuid) {
|
||||
break;
|
||||
case CENTER: {
|
||||
std::vector<std::pair<ChunkPos, PendingChunk*>> tasks;
|
||||
for (auto& [pos, task] : new_chunks) {
|
||||
for (auto& [pos, task] : m_new_chunks) {
|
||||
if (!task.future.valid()) {
|
||||
tasks.emplace_back(pos, &task);
|
||||
}
|
||||
@@ -177,7 +180,7 @@ void ServerWorld::poll_finished_chunks() {
|
||||
m_new_finished_chunk.clear();
|
||||
std::lock_guard lock(m_new_chunk_mutex);
|
||||
std::erase_if(
|
||||
new_chunks, [&](std::pair<const ChunkPos, PendingChunk>& pair) {
|
||||
m_new_chunks, [&](std::pair<const ChunkPos, PendingChunk>& pair) {
|
||||
auto& pending = pair.second;
|
||||
if (!pending.future.valid()) {
|
||||
return false;
|
||||
@@ -200,9 +203,9 @@ void ServerWorld::start_gen_thread() {
|
||||
while (!token.stop_requested()) {
|
||||
std::unique_lock<std::mutex> lk(m_need_gen_queue_mutex);
|
||||
|
||||
m_gen_cv.wait(lk, [this](std::stop_token token) {
|
||||
m_gen_cv.wait(lk, token, [this]() {
|
||||
return m_need_gen_chunk.load() || !m_gen_running ||
|
||||
!m_need_gen_queue.empty() || token.stop_requested();
|
||||
!m_need_gen_queue.empty();
|
||||
});
|
||||
if (!m_gen_running) {
|
||||
break;
|
||||
@@ -330,10 +333,9 @@ void ServerWorld::hot_reload() {
|
||||
}
|
||||
|
||||
void ServerWorld::rebuild_world() {
|
||||
if (m_is_rebuilding) {
|
||||
if (m_is_rebuilding.exchange(true)) {
|
||||
return;
|
||||
}
|
||||
m_is_rebuilding = true;
|
||||
stop_gen_thread();
|
||||
stop_thread_pool();
|
||||
m_cave_carcer.reload(ChunkGenerator::seed());
|
||||
@@ -343,22 +345,40 @@ void ServerWorld::rebuild_world() {
|
||||
m_chunks.clear();
|
||||
m_new_finished_chunk.clear();
|
||||
}
|
||||
{
|
||||
std::lock_guard lock(m_new_chunk_mutex);
|
||||
m_new_chunks.clear();
|
||||
}
|
||||
m_could_gen = true;
|
||||
ChunkGenerator::reload();
|
||||
start_thread_pool();
|
||||
start_gen_thread();
|
||||
need_gen(std::nullopt);
|
||||
|
||||
m_is_rebuilding = false;
|
||||
}
|
||||
|
||||
void ServerWorld::update() { poll_finished_chunks(); }
|
||||
void ServerWorld::update() {
|
||||
poll_finished_chunks();
|
||||
{
|
||||
std::lock_guard lk(m_chunks_mutex);
|
||||
bool consumed = false;
|
||||
|
||||
for (auto& x : m_new_finished_chunk) {
|
||||
m_chunks.emplace(x.first, std::move(x.second));
|
||||
consumed = true;
|
||||
}
|
||||
if (consumed) {
|
||||
m_could_gen = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void ServerWorld::sync_player_pos(const std::string& uuid, float x, float y,
|
||||
float z) {
|
||||
std::lock_guard lock(m_player_mutex);
|
||||
auto it = m_players.find(uuid);
|
||||
if (it == m_players.end()) {
|
||||
Logger::warn("Player {} is not in this Server", it->first);
|
||||
Logger::warn("Player {} is not in this Server", uuid);
|
||||
return;
|
||||
}
|
||||
it->second.update_pos(x, y, z);
|
||||
@@ -367,31 +387,36 @@ void ServerWorld::sync_player_pos(const std::string& uuid, float x, float y,
|
||||
void ServerWorld::handle_player_login(const std::string& name,
|
||||
std::shared_ptr<Session> session) {
|
||||
std::string uuid = generate_uuid();
|
||||
player_join(name, uuid);
|
||||
m_player_session.emplace(uuid, session);
|
||||
Logger::info("Player {} (uuid {}) join the world", name, uuid);
|
||||
{
|
||||
std::lock_guard lock(m_player_mutex);
|
||||
m_players.emplace(std::piecewise_construct,
|
||||
std::forward_as_tuple(std::string(uuid)),
|
||||
std::forward_as_tuple(name, uuid, *this, session));
|
||||
}
|
||||
m_uuid_to_name.emplace(uuid, name);
|
||||
LoginRsp rsp;
|
||||
rsp.set_success(true);
|
||||
rsp.set_uuid(uuid);
|
||||
session->send(make_packet(name));
|
||||
session->send(make_packet(rsp));
|
||||
}
|
||||
|
||||
void ServerWorld::player_join(std::string_view name, std::string_view uuid) {
|
||||
Logger::info("Player {} (uuid {}) join the world", name, uuid);
|
||||
m_players.emplace(std::piecewise_construct,
|
||||
std::forward_as_tuple(std::string(uuid)),
|
||||
std::forward_as_tuple(name, uuid, *this));
|
||||
}
|
||||
void ServerWorld::handle_player_exit(const std::string& uuid) {
|
||||
{
|
||||
std::lock_guard lock(m_player_mutex);
|
||||
auto it = m_players.find(uuid);
|
||||
if (it != m_players.end()) {
|
||||
|
||||
void ServerWorld::player_exit(const std::string& name) {
|
||||
auto it = m_players.find(name);
|
||||
if (it == m_players.end()) {
|
||||
Logger::error("Player {} isn't in Server", it->first);
|
||||
m_players.erase(it);
|
||||
} else {
|
||||
Logger::error("Player {} isn't in Server", uuid);
|
||||
}
|
||||
m_players.unsafe_erase(it);
|
||||
}
|
||||
m_uuid_to_name.erase(uuid);
|
||||
}
|
||||
|
||||
glm::vec3 ServerWorld::get_player_pos(const std::string& uuid) const {
|
||||
std::shared_lock lock(m_player_mutex);
|
||||
auto it = m_players.find(uuid);
|
||||
if (it == m_players.end()) {
|
||||
Logger::error("Can't find player uuid {}", uuid);
|
||||
@@ -437,9 +462,10 @@ void ServerWorld::handle_chunk_req(const std::string& uuid, ChunkPos pos) {
|
||||
}
|
||||
std::shared_ptr<Session> s;
|
||||
{
|
||||
session_cacc cacc;
|
||||
if (m_player_session.find(cacc, uuid)) {
|
||||
s = cacc->second;
|
||||
std::shared_lock lock(m_player_mutex);
|
||||
auto it = m_players.find(uuid);
|
||||
if (it != m_players.end()) {
|
||||
s = it->second.get_session();
|
||||
}
|
||||
}
|
||||
if (!s) {
|
||||
@@ -462,12 +488,15 @@ void ServerWorld::handle_block_change(const BlockChangeReq& req) {
|
||||
pos->set_y(y);
|
||||
pos->set_z(z);
|
||||
rsp.set_block(req.block());
|
||||
|
||||
for (auto& [uuid, session] : m_player_session) {
|
||||
{
|
||||
std::shared_lock lock(m_player_mutex);
|
||||
for (auto& [uuid, player] : m_players) {
|
||||
auto session = player.get_session();
|
||||
if (session) {
|
||||
session->send(make_packet(rsp));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
int ServerWorld::rendering_distance() const {
|
||||
|
||||
@@ -88,7 +88,15 @@ asio::awaitable<void> Session::read_loop() {
|
||||
}
|
||||
}
|
||||
} catch (const asio::system_error& e) {
|
||||
Logger::warn("Catch Asio Error {}", e.what());
|
||||
auto ec = e.code();
|
||||
|
||||
if (ec == asio::error::eof || ec == asio::error::operation_aborted) {
|
||||
|
||||
Logger::info("Client disconnected");
|
||||
} else {
|
||||
Logger::warn("Asio Error {}", e.what());
|
||||
}
|
||||
|
||||
close();
|
||||
} catch (...) {
|
||||
Logger::error("Unknow Error");
|
||||
|
||||
Reference in New Issue
Block a user