常规更新

This commit is contained in:
2026-07-24 10:54:45 +08:00
parent 0f4623cbef
commit e8bcc7ecfd
117 changed files with 12097 additions and 14121 deletions
+9 -20
View File
@@ -1,40 +1,29 @@
#pragma once
#include <system_error>
#include "Core/Base/RingBuffer.hpp"
#include "Socket.h"
#include "global.h"
#include <system_error>
namespace Psc::asio_socket {
#ifndef PSC_ASIO_STREAM_RING_BUFFER
#define PSC_ASIO_STREAM_RING_BUFFER StreamRingBuffer_ST
#endif
#ifndef PSC_ASIO_PACKET_RING_BUFFER
#define PSC_ASIO_PACKET_RING_BUFFER RingBuffer_ST
#endif
using ASIO_StreamRingBuffer = PSC_ASIO_STREAM_RING_BUFFER;
using ASIO_PacketRingBuffer = PSC_ASIO_PACKET_RING_BUFFER;
inline Socket_FD socket_id(const void *ptr) {
return static_cast<Socket_FD>(reinterpret_cast<std::uintptr_t>(ptr));
inline Socket_FD socket_id(const void* ptr) {
return static_cast<Socket_FD>(reinterpret_cast<std::uintptr_t>(ptr));
}
inline Sockaddr_In
endpoint_to_sockaddr(const asio::ip::tcp::endpoint &endpoint) {
return {endpoint.address().to_string(), endpoint.port()};
endpoint_to_sockaddr(const asio::ip::tcp::endpoint& endpoint) {
return {endpoint.address().to_string(), endpoint.port()};
}
inline Sockaddr_In
endpoint_to_sockaddr(const asio::ip::udp::endpoint &endpoint) {
return {endpoint.address().to_string(), endpoint.port()};
endpoint_to_sockaddr(const asio::ip::udp::endpoint& endpoint) {
return {endpoint.address().to_string(), endpoint.port()};
}
inline bool would_block(const asio::error_code &ec) {
return ec == asio::error::would_block || ec == asio::error::try_again;
inline bool would_block(const asio::error_code& ec) {
return ec == asio::error::would_block || ec == asio::error::try_again;
}
} // namespace Psc::asio_socket
-1
View File
@@ -1,5 +1,4 @@
#pragma once
#if defined(_WIN32)
#ifndef NOMINMAX
#define NOMINMAX
-4
View File
@@ -1,8 +1,4 @@
#include "Socket.h"
namespace Psc::asio_socket {
BaseLogger* socket_logger = nullptr;
} // namespace Psc::asio_socket
+8 -26
View File
@@ -1,69 +1,54 @@
#pragma once
#include <string_view>
#include "../Base/JSON.h"
#include "../system/export.h"
#include "Core/spdlog/export.h"
#include <cstdint>
#include <optional>
#include <string>
#include <string_view>
#include "../Base/JSON.h"
#include "../system/export.h"
#include "Core/spdlog/export.h"
namespace Psc::asio_socket {
#if _WIN32
using Socket_FD = unsigned long long;
#else
using Socket_FD = int;
#endif
extern BaseLogger* socket_logger;
struct Sockaddr_In {
std::string ip;
std::uint32_t port{};
Sockaddr_In() = default;
Sockaddr_In(std::string_view ip, std::uint32_t port) : ip(ip), port(port) {}
Sockaddr_In(std::string_view ip, std::uint32_t port) : ip(ip), port(port) {
}
[[nodiscard]] std::string to_string() const {
return ip + ":" + std::to_string(port);
}
bool operator<(const Sockaddr_In& other) const {
if (ip != other.ip) return ip < other.ip;
if (ip != other.ip)
return ip < other.ip;
return port < other.port;
}
bool operator==(const Sockaddr_In& other) const noexcept {
return ip == other.ip && port == other.port;
}
bool operator!=(const Sockaddr_In& other) const noexcept {
return !(*this == other);
}
};
struct Accept_Info {
Sockaddr_In sockaddr;
Socket_FD fd{};
[[nodiscard]] std::string to_string() const {
return "fd:" + std::to_string(fd) + ":[" + sockaddr.to_string() + "]";
}
};
class Socket_Base {
public:
virtual ~Socket_Base() = default;
std::optional<Sockaddr_In> address;
Socket_FD socket_fd = static_cast<Socket_FD>(-1);
virtual std::string to_string() {
return address ? address->to_string() : std::string{};
}
[[nodiscard]] JSON to_Json() const {
JSON ret = JSON::object();
ret.append({"fd", static_cast<std::uint64_t>(socket_fd)});
@@ -73,11 +58,8 @@ public:
}
return ret;
}
virtual void close() {
socket_fd = static_cast<Socket_FD>(-1);
}
};
} // namespace Psc::asio_socket
+136 -135
View File
@@ -1,10 +1,4 @@
#pragma once
#include <string_view>
#include "ASIO_Utils.h"
#include "TCP_Client.h"
#include "TCP_Server.h"
#include "UDP_Client.h"
#include "UDP_Server.h"
#include <array>
#include <asio/awaitable.hpp>
#include <asio/use_awaitable.hpp>
@@ -12,180 +6,187 @@
#include <cstdint>
#include <memory>
#include <string>
#include <string_view>
#include <system_error>
#include <utility>
#include <vector>
#include "ASIO_Utils.h"
#include "TCP_Client.h"
#include "TCP_Server.h"
#include "UDP_Client.h"
#include "UDP_Server.h"
namespace Psc::asio_socket {
class TCP_Client_Coro : public TCP_Client {
public:
[[nodiscard]] asio::awaitable<bool> connect_coro() {
auto started = start_connect();
tick();
co_return started;
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<std::string> read_coro() {
co_return read();
}
[[nodiscard]] asio::awaitable<void> send_coro(std::string_view data) {
send(data);
co_return;
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
[[nodiscard]] asio::awaitable<bool> connect_coro() {
auto started = start_connect();
tick();
co_return started;
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<std::string> read_coro() {
co_return read();
}
[[nodiscard]] asio::awaitable<void> send_coro(std::string_view data) {
send(data);
co_return;
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
};
class TCP_Server_Coro : public TCP_Server {
public:
[[nodiscard]] asio::awaitable<bool> listen_coro(std::string_view ip,
std::uint32_t port) {
co_return listen(std::move(ip), port);
}
[[nodiscard]] asio::awaitable<bool> listen_coro(const Sockaddr_In& address) {
co_return listen(address);
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<void> flush_clients_coro() {
flush_clients();
co_return;
}
[[nodiscard]] asio::awaitable<void> write_to_all_clients_coro(std::string_view data) {
write_to_all_clients(data);
co_return;
}
[[nodiscard]] asio::awaitable<std::vector<Read_Info>> read_from_all_clients_coro() {
co_return read_from_all_clients();
}
[[nodiscard]] asio::awaitable<std::shared_ptr<TCP_Connect>> accept_coro() const {
co_return accept();
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
[[nodiscard]] asio::awaitable<bool> listen_coro(std::string_view ip,
std::uint32_t port) {
co_return listen(std::move(ip), port);
}
[[nodiscard]] asio::awaitable<bool> listen_coro(const Sockaddr_In& address) {
co_return listen(address);
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<void> flush_clients_coro() {
flush_clients();
co_return;
}
[[nodiscard]] asio::awaitable<void> write_to_all_clients_coro(std::string_view data) {
write_to_all_clients(data);
co_return;
}
[[nodiscard]] asio::awaitable<std::vector<Read_Info>> read_from_all_clients_coro() {
co_return read_from_all_clients();
}
[[nodiscard]] asio::awaitable<std::shared_ptr<TCP_Connect>> accept_coro() const {
co_return accept();
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
};
class UDP_Client_Coro : public UDP_Client {
public:
[[nodiscard]] asio::awaitable<bool> connect_coro() {
co_return connect();
}
[[nodiscard]] asio::awaitable<std::string> read_coro() {
co_return read();
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<void> send_coro(std::string_view data) {
send(data);
co_return;
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
[[nodiscard]] asio::awaitable<bool> connect_coro() {
co_return connect();
}
[[nodiscard]] asio::awaitable<std::string> read_coro() {
co_return read();
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<void> send_coro(std::string_view data) {
send(data);
co_return;
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
};
class UDP_Server_Coro : public UDP_Server {
public:
[[nodiscard]] asio::awaitable<bool> bind_coro() {
co_return bind();
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<std::vector<Read_Info>> read_coro() {
co_return read();
}
[[nodiscard]] asio::awaitable<void> reply_last_peer_coro(std::string_view data) {
reply_last_peer(data);
co_return;
}
[[nodiscard]] asio::awaitable<void> send_to_coro(std::string_view ip, uint16_t port, std::string_view data) {
send_to(ip, port, data);
co_return;
}
[[nodiscard]] asio::awaitable<void> write_to_all_clients_coro(std::string_view msg) {
write_to_all_clients(msg);
co_return;
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
[[nodiscard]] asio::awaitable<bool> bind_coro() {
co_return bind();
}
[[nodiscard]] asio::awaitable<void> tick_coro() {
tick();
co_return;
}
[[nodiscard]] asio::awaitable<std::vector<Read_Info>> read_coro() {
co_return read();
}
[[nodiscard]] asio::awaitable<void> reply_last_peer_coro(std::string_view data) {
reply_last_peer(data);
co_return;
}
[[nodiscard]] asio::awaitable<void> send_to_coro(std::string_view ip, uint16_t port, std::string_view data) {
send_to(ip, port, data);
co_return;
}
[[nodiscard]] asio::awaitable<void> write_to_all_clients_coro(std::string_view msg) {
write_to_all_clients(msg);
co_return;
}
[[nodiscard]] asio::awaitable<void> close_coro() {
close();
co_return;
}
};
} // namespace Psc::asio_socket
namespace Psc::asio_socket::coro {
struct TCP_Read_Result {
std::vector<char> data;
std::vector<char> data;
};
struct UDP_Read_Result {
Sockaddr_In remote;
std::vector<char> data;
Sockaddr_In remote;
std::vector<char> data;
};
inline void throw_if_error(const asio::error_code& ec) {
if (ec) {
throw std::system_error(ec);
}
if (ec) {
throw std::system_error(ec);
}
}
inline asio::awaitable<
std::shared_ptr<asio::ip::tcp::resolver::results_type>> tcp_resolve(asio::ip::tcp::resolver& resolver, std::string_view host,
std::uint16_t port) {
auto endpoints =
co_await resolver.async_resolve(host, std::to_string(port), asio::use_awaitable);
co_return std::make_shared<asio::ip::tcp::resolver::results_type>(std::move(endpoints));
std::shared_ptr<asio::ip::tcp::resolver::results_type>>
tcp_resolve(asio::ip::tcp::resolver& resolver, std::string_view host,
std::uint16_t port) {
auto endpoints =
co_await resolver.async_resolve(host, std::to_string(port), asio::use_awaitable);
co_return std::make_shared<asio::ip::tcp::resolver::results_type>(std::move(endpoints));
}
inline asio::awaitable<asio::ip::tcp::endpoint> tcp_connect(asio::ip::tcp::socket& socket,
const asio::ip::tcp::resolver::results_type& endpoints) {
co_return co_await asio::async_connect(socket, endpoints, asio::use_awaitable);
co_return co_await asio::async_connect(socket, endpoints, asio::use_awaitable);
}
inline asio::awaitable<asio::ip::tcp::endpoint> tcp_connect(asio::ip::tcp::socket& socket, asio::ip::tcp::resolver& resolver,
std::string_view host, std::uint16_t port) {
auto endpoints = co_await tcp_resolve(resolver, std::move(host), port);
co_return co_await tcp_connect(socket, *endpoints);
auto endpoints = co_await tcp_resolve(resolver, std::move(host), port);
co_return co_await tcp_connect(socket, *endpoints);
}
inline asio::awaitable<std::shared_ptr<asio::ip::tcp::socket>> tcp_accept(asio::ip::tcp::acceptor& acceptor) {
auto socket =
std::make_shared<asio::ip::tcp::socket>(acceptor.get_executor());
co_await acceptor.async_accept(*socket, asio::use_awaitable);
co_return socket;
auto socket =
std::make_shared<asio::ip::tcp::socket>(acceptor.get_executor());
co_await acceptor.async_accept(*socket, asio::use_awaitable);
co_return socket;
}
inline asio::awaitable<TCP_Read_Result> tcp_read_some(asio::ip::tcp::socket& socket, std::size_t max_size = 16 * 1024) {
std::vector<char> buffer(max_size);
auto size = co_await socket.async_read_some(asio::buffer(buffer), asio::use_awaitable);
buffer.resize(size);
co_return TCP_Read_Result{std::move(buffer)};
std::vector<char> buffer(max_size);
auto size = co_await socket.async_read_some(asio::buffer(buffer), asio::use_awaitable);
buffer.resize(size);
co_return TCP_Read_Result{std::move(buffer)};
}
inline asio::awaitable<std::size_t> tcp_write(asio::ip::tcp::socket& socket, std::string_view data) {
co_return co_await asio::async_write(socket, asio::buffer(data), asio::use_awaitable);
co_return co_await asio::async_write(socket, asio::buffer(data), asio::use_awaitable);
}
inline asio::awaitable<UDP_Read_Result> udp_receive_from(asio::ip::udp::socket& socket,
std::size_t max_size = 16 * 1024) {
std::vector<char> buffer(max_size);
asio::ip::udp::endpoint remote;
auto size = co_await socket.async_receive_from(asio::buffer(buffer), remote, asio::use_awaitable);
buffer.resize(size);
co_return UDP_Read_Result{endpoint_to_sockaddr(remote), std::move(buffer)};
std::vector<char> buffer(max_size);
asio::ip::udp::endpoint remote;
auto size = co_await socket.async_receive_from(asio::buffer(buffer), remote, asio::use_awaitable);
buffer.resize(size);
co_return UDP_Read_Result{endpoint_to_sockaddr(remote), std::move(buffer)};
}
inline asio::awaitable<std::size_t> udp_send_to(asio::ip::udp::socket& socket, std::string_view data,
asio::ip::udp::endpoint remote) {
co_return co_await socket.async_send_to(asio::buffer(data), remote, asio::use_awaitable);
co_return co_await socket.async_send_to(asio::buffer(data), remote, asio::use_awaitable);
}
inline asio::awaitable<std::size_t> udp_send_to(asio::ip::udp::socket& socket, std::string_view data,
const Sockaddr_In& remote) {
asio::error_code ec;
auto address = asio::ip::make_address(remote.ip, ec);
throw_if_error(ec);
co_return co_await udp_send_to(
socket, std::move(data),
asio::ip::udp::endpoint(address,
static_cast<unsigned short>(remote.port)));
asio::error_code ec;
auto address = asio::ip::make_address(remote.ip, ec);
throw_if_error(ec);
co_return co_await udp_send_to(
socket, std::move(data),
asio::ip::udp::endpoint(address,
static_cast<unsigned short>(remote.port)));
}
} // namespace Psc::asio_socket::coro
+144 -167
View File
@@ -1,206 +1,183 @@
#include "TCP_Client.h"
#include <array>
#include <string_view>
namespace Psc::asio_socket {
TCP_Client::~TCP_Client() { close(); }
TCP_Client::~TCP_Client() {
close();
}
void TCP_Client::set_state(State s) {
const auto old = state.exchange(s);
if (old != s && on_state_change)
on_state_change(old, s);
const auto old = state.exchange(s);
if (old != s && on_state_change)
on_state_change(old, s);
}
void TCP_Client::create() {
send_buffer.init(max_buffer_size);
recv_buffer.init(max_buffer_size);
connect_pending = false;
read_pending = false;
write_pending = false;
send_storage.clear();
resolver = std::make_unique<asio::ip::tcp::resolver>(io_context);
socket = std::make_unique<asio::ip::tcp::socket>(io_context);
socket_fd = socket_id(socket.get());
set_state(Reconnecting);
send_buffer.init(max_buffer_size);
recv_buffer.init(max_buffer_size);
connect_pending = false;
read_pending = false;
write_pending = false;
send_storage.clear();
resolver = std::make_unique<asio::ip::tcp::resolver>(io_context);
socket = std::make_unique<asio::ip::tcp::socket>(io_context);
socket_fd = socket_id(socket.get());
set_state(Reconnecting);
}
void TCP_Client::close() {
connect_pending = false;
read_pending = false;
write_pending = false;
if (resolver) {
asio::error_code ec;
resolver->cancel();
}
if (socket) {
asio::error_code ec;
socket->cancel(ec);
socket->shutdown(asio::ip::tcp::socket::shutdown_both, ec);
socket->close(ec);
}
socket_fd = static_cast<Socket_FD>(-1);
set_state(Not_Created);
connect_pending = false;
read_pending = false;
write_pending = false;
if (resolver) {
asio::error_code ec;
resolver->cancel();
}
if (socket) {
asio::error_code ec;
socket->cancel(ec);
socket->shutdown(asio::ip::tcp::socket::shutdown_both, ec);
socket->close(ec);
}
socket_fd = static_cast<Socket_FD>(-1);
set_state(Not_Created);
}
void TCP_Client::schedule_recreate() {
connect_pending = false;
read_pending = false;
write_pending = false;
if (resolver)
resolver->cancel();
if (socket) {
asio::error_code ec;
socket->cancel(ec);
socket->close(ec);
}
next_reconnect_tp = std::chrono::steady_clock::now() +
std::chrono::milliseconds(reconnect_backoff_ms);
set_state(Reconnecting);
connect_pending = false;
read_pending = false;
write_pending = false;
if (resolver)
resolver->cancel();
if (socket) {
asio::error_code ec;
socket->cancel(ec);
socket->close(ec);
}
next_reconnect_tp = std::chrono::steady_clock::now() +
std::chrono::milliseconds(reconnect_backoff_ms);
set_state(Reconnecting);
}
bool TCP_Client::start_connect() {
if (connect_pending || state == Connecting || state == Connected)
return true;
if (dest_address.ip.empty() || dest_address.port == 0) {
set_state(Not_Set_Field);
return false;
}
io_context.restart();
resolver = std::make_unique<asio::ip::tcp::resolver>(io_context);
socket = std::make_unique<asio::ip::tcp::socket>(io_context);
connect_pending = true;
set_state(Connecting);
resolver->async_resolve(
dest_address.ip, std::to_string(dest_address.port),
[this](const asio::error_code &ec,
asio::ip::tcp::resolver::results_type endpoints) {
if (connect_pending || state == Connecting || state == Connected)
return true;
if (dest_address.ip.empty() || dest_address.port == 0) {
set_state(Not_Set_Field);
return false;
}
io_context.restart();
resolver = std::make_unique<asio::ip::tcp::resolver>(io_context);
socket = std::make_unique<asio::ip::tcp::socket>(io_context);
connect_pending = true;
set_state(Connecting);
resolver->async_resolve(
dest_address.ip, std::to_string(dest_address.port),
[this](const asio::error_code& ec,
asio::ip::tcp::resolver::results_type endpoints) {
if (state != Connecting || !socket)
return;
return;
if (ec) {
schedule_recreate();
return;
schedule_recreate();
return;
}
asio::async_connect(*socket, endpoints,
[this](const asio::error_code &connect_ec,
const asio::ip::tcp::endpoint &) {
connect_pending = false;
if (state != Connecting || !socket)
return;
if (connect_ec) {
schedule_recreate();
return;
}
socket_fd = socket_id(socket.get());
address = dest_address;
set_state(Connected);
start_read();
start_write();
[this](const asio::error_code& connect_ec,
const asio::ip::tcp::endpoint&) {
connect_pending = false;
if (state != Connecting || !socket)
return;
if (connect_ec) {
schedule_recreate();
return;
}
socket_fd = socket_id(socket.get());
address = dest_address;
set_state(Connected);
start_read();
start_write();
});
});
return true;
});
return true;
}
void TCP_Client::tick() {
if (state == Reconnecting &&
std::chrono::steady_clock::now() >= next_reconnect_tp) {
start_connect();
}
if (state == Connected) {
start_read();
start_write();
}
io_context.restart();
io_context.poll();
if (state == Connected) {
start_read();
start_write();
}
if (state == Reconnecting &&
std::chrono::steady_clock::now() >= next_reconnect_tp) {
start_connect();
}
if (state == Connected) {
start_read();
start_write();
}
io_context.restart();
io_context.poll();
if (state == Connected) {
start_read();
start_write();
}
}
void TCP_Client::start_read() {
if (read_pending || state != Connected || !socket || !socket->is_open())
return;
read_pending = true;
socket->async_read_some(asio::buffer(recv_storage),
[this](const asio::error_code &ec, std::size_t n) {
read_pending = false;
if (state != Connected)
return;
if (ec == asio::error::operation_aborted)
return;
if (ec == asio::error::eof ||
ec == asio::error::connection_reset || ec) {
schedule_recreate();
return;
}
recv_buffer.write_best_effort(recv_storage.data(),
n);
start_read();
});
if (read_pending || state != Connected || !socket || !socket->is_open())
return;
read_pending = true;
socket->async_read_some(asio::buffer(recv_storage),
[this](const asio::error_code& ec, std::size_t n) {
read_pending = false;
if (state != Connected)
return;
if (ec == asio::error::operation_aborted)
return;
if (ec == asio::error::eof ||
ec == asio::error::connection_reset || ec) {
schedule_recreate();
return;
}
recv_buffer.write_best_effort(recv_storage.data(),
n);
start_read();
});
}
void TCP_Client::start_write() {
if (write_pending || state != Connected || !socket || !socket->is_open())
return;
send_storage.resize(16 * 1024);
auto size =
send_buffer.peek_best_effort(send_storage.data(), send_storage.size());
if (size == 0)
return;
send_storage.resize(size);
write_pending = true;
socket->async_write_some(
asio::buffer(send_storage.data(), send_storage.size()),
[this](const asio::error_code &ec, std::size_t sent) {
if (write_pending || state != Connected || !socket || !socket->is_open())
return;
send_storage.resize(16 * 1024);
auto size =
send_buffer.peek_best_effort(send_storage.data(), send_storage.size());
if (size == 0)
return;
send_storage.resize(size);
write_pending = true;
socket->async_write_some(
asio::buffer(send_storage.data(), send_storage.size()),
[this](const asio::error_code& ec, std::size_t sent) {
write_pending = false;
if (state != Connected)
return;
return;
if (ec == asio::error::operation_aborted)
return;
return;
if (ec) {
schedule_recreate();
return;
schedule_recreate();
return;
}
send_buffer.skip(sent);
if (sent != 0)
start_write();
});
start_write();
});
}
std::string TCP_Client::read() {
std::string ret;
std::array<char, 16 * 1024> buffer{};
for (;;) {
auto n = recv_buffer.read_best_effort(buffer.data(), buffer.size());
if (n == 0)
break;
ret.append(buffer.data(), n);
}
return ret;
std::string ret;
std::array<char, 16 * 1024> buffer{};
for (;;) {
auto n = recv_buffer.read_best_effort(buffer.data(), buffer.size());
if (n == 0)
break;
ret.append(buffer.data(), n);
}
return ret;
}
void TCP_Client::send(std::string_view data) {
if (data.empty())
return;
send_buffer.write_best_effort(data.data(), data.size());
if (state == Connected)
start_write();
if (data.empty())
return;
send_buffer.write_best_effort(data.data(), data.size());
if (state == Connected)
start_write();
}
std::string TCP_Client::to_string() {
return "TCP_Client:[" + dest_address.to_string() + "]";
return "TCP_Client:[" + dest_address.to_string() + "]";
}
} // namespace Psc::asio_socket
+45 -52
View File
@@ -1,65 +1,58 @@
#pragma once
#include <string_view>
#include "ASIO_Utils.h"
#include <array>
#include <atomic>
#include <chrono>
#include <functional>
#include <memory>
#include <string_view>
#include <vector>
#include "ASIO_Utils.h"
namespace Psc::asio_socket {
class TCP_Client : public Socket_Base {
public:
~TCP_Client() override;
void set_buffer_size(size_t size) { max_buffer_size = size; }
void set_dest_address(const Sockaddr_In &a) { dest_address = a; }
void create();
void tick();
std::string read();
void send(std::string_view data);
std::string to_string() override;
enum State {
Not_Created,
Not_Set_Field,
Reconnecting,
Connecting,
Wait_Check,
Connected
};
void close() override;
std::atomic<State> state{Not_Created};
void set_state(State s);
std::function<void(State old_state, State new_state)> on_state_change =
nullptr;
~TCP_Client() override;
void set_buffer_size(size_t size) {
max_buffer_size = size;
}
void set_dest_address(const Sockaddr_In& a) {
dest_address = a;
}
void create();
void tick();
std::string read();
void send(std::string_view data);
std::string to_string() override;
enum State {
Not_Created,
Not_Set_Field,
Reconnecting,
Connecting,
Wait_Check,
Connected
};
void close() override;
std::atomic<State> state{Not_Created};
void set_state(State s);
std::function<void(State old_state, State new_state)> on_state_change =
nullptr;
protected:
size_t max_buffer_size = 1024 * 1024;
ASIO_StreamRingBuffer send_buffer;
ASIO_StreamRingBuffer recv_buffer;
Sockaddr_In dest_address{};
int reconnect_backoff_ms = 200;
std::chrono::steady_clock::time_point next_reconnect_tp{};
asio::io_context io_context;
std::unique_ptr<asio::ip::tcp::resolver> resolver;
std::unique_ptr<asio::ip::tcp::socket> socket;
std::array<char, 16 * 1024> recv_storage{};
std::vector<char> send_storage;
bool connect_pending = false;
bool read_pending = false;
bool write_pending = false;
void schedule_recreate();
bool start_connect();
void start_read();
void start_write();
size_t max_buffer_size = 1024 * 1024;
ASIO_StreamRingBuffer send_buffer;
ASIO_StreamRingBuffer recv_buffer;
Sockaddr_In dest_address{};
int reconnect_backoff_ms = 200;
std::chrono::steady_clock::time_point next_reconnect_tp{};
asio::io_context io_context;
std::unique_ptr<asio::ip::tcp::resolver> resolver;
std::unique_ptr<asio::ip::tcp::socket> socket;
std::array<char, 16 * 1024> recv_storage{};
std::vector<char> send_storage;
bool connect_pending = false;
bool read_pending = false;
bool write_pending = false;
void schedule_recreate();
bool start_connect();
void start_read();
void start_write();
};
} // namespace Psc::asio_socket
+230 -268
View File
@@ -1,329 +1,291 @@
#include "TCP_Server.h"
#include <array>
#include <string_view>
namespace Psc::asio_socket {
TCP_Server::~TCP_Server() { close(); }
TCP_Server::~TCP_Server() {
close();
}
void TCP_Server::tick() {
if (state == Working) {
start_accept();
flush_clients();
}
io_context.restart();
io_context.poll();
if (state == Working) {
flush_clients();
cleanup_closed_clients();
}
if (state == Working) {
start_accept();
flush_clients();
}
io_context.restart();
io_context.poll();
if (state == Working) {
flush_clients();
cleanup_closed_clients();
}
}
TCP_Server &TCP_Server::set_tcp_no_delay(bool value) {
no_delay = value;
return *this;
TCP_Server& TCP_Server::set_tcp_no_delay(bool value) {
no_delay = value;
return *this;
}
TCP_Server &TCP_Server::set_connect_system_buffer_size(size_t size) {
connect_system_buffer_size = size;
return *this;
TCP_Server& TCP_Server::set_connect_system_buffer_size(size_t size) {
connect_system_buffer_size = size;
return *this;
}
TCP_Server &TCP_Server::set_connect_user_buffer_size(size_t size) {
connect_user_buffer_size = size;
return *this;
TCP_Server& TCP_Server::set_connect_user_buffer_size(size_t size) {
connect_user_buffer_size = size;
return *this;
}
TCP_Server &TCP_Server::set_recv_system_buffer_size(size_t size) {
recv_system_buffer_size = size;
return *this;
TCP_Server& TCP_Server::set_recv_system_buffer_size(size_t size) {
recv_system_buffer_size = size;
return *this;
}
TCP_Server &TCP_Server::create() {
acceptor = std::make_unique<asio::ip::tcp::acceptor>(io_context);
socket_fd = socket_id(acceptor.get());
state = Not_bind;
return *this;
TCP_Server& TCP_Server::create() {
acceptor = std::make_unique<asio::ip::tcp::acceptor>(io_context);
socket_fd = socket_id(acceptor.get());
state = Not_bind;
return *this;
}
bool TCP_Server::listen(std::string_view ip, std::uint32_t port) {
return listen(Sockaddr_In(ip, port));
return listen(Sockaddr_In(ip, port));
}
bool TCP_Server::listen(const Sockaddr_In &addr) {
if (!acceptor)
create();
if (addr.ip.empty() || addr.port == 0) {
state = Not_Set_Field;
return false;
}
asio::error_code ec;
const auto address_value = asio::ip::make_address(addr.ip, ec);
if (ec) {
state = Not_bind;
return false;
}
asio::ip::tcp::endpoint endpoint(address_value,
static_cast<unsigned short>(addr.port));
acceptor->open(endpoint.protocol(), ec);
if (ec)
return false;
acceptor->set_option(asio::socket_base::reuse_address(true), ec);
acceptor->bind(endpoint, ec);
if (ec) {
state = Not_bind;
return false;
}
acceptor->listen(asio::socket_base::max_listen_connections, ec);
if (ec) {
state = Not_Listen;
return false;
}
address = addr;
state = Working;
start_accept();
return true;
bool TCP_Server::listen(const Sockaddr_In& addr) {
if (!acceptor)
create();
if (addr.ip.empty() || addr.port == 0) {
state = Not_Set_Field;
return false;
}
asio::error_code ec;
const auto address_value = asio::ip::make_address(addr.ip, ec);
if (ec) {
state = Not_bind;
return false;
}
asio::ip::tcp::endpoint endpoint(address_value,
static_cast<unsigned short>(addr.port));
acceptor->open(endpoint.protocol(), ec);
if (ec)
return false;
acceptor->set_option(asio::socket_base::reuse_address(true), ec);
acceptor->bind(endpoint, ec);
if (ec) {
state = Not_bind;
return false;
}
acceptor->listen(asio::socket_base::max_listen_connections, ec);
if (ec) {
state = Not_Listen;
return false;
}
address = addr;
state = Working;
start_accept();
return true;
}
void TCP_Server::close() {
accept_pending = false;
if (pending_accept_socket) {
asio::error_code ec;
pending_accept_socket->cancel(ec);
pending_accept_socket->close(ec);
}
pending_accept_socket.reset();
for (auto &[_, conn] : tcp_clients) {
close_client(conn);
}
tcp_clients.clear();
accepted_clients.clear();
if (acceptor) {
asio::error_code ec;
acceptor->cancel(ec);
acceptor->close(ec);
}
socket_fd = static_cast<Socket_FD>(-1);
state = Not_Created;
accept_pending = false;
if (pending_accept_socket) {
asio::error_code ec;
pending_accept_socket->cancel(ec);
pending_accept_socket->close(ec);
}
pending_accept_socket.reset();
for (auto& [_, conn] : tcp_clients) {
close_client(conn);
}
tcp_clients.clear();
accepted_clients.clear();
if (acceptor) {
asio::error_code ec;
acceptor->cancel(ec);
acceptor->close(ec);
}
socket_fd = static_cast<Socket_FD>(-1);
state = Not_Created;
}
std::shared_ptr<TCP_Connect> TCP_Server::accept() const {
while (!accepted_clients.empty()) {
auto conn = accepted_clients.front().lock();
accepted_clients.erase(accepted_clients.begin());
if (conn && !conn->closing)
return conn;
}
return nullptr;
while (!accepted_clients.empty()) {
auto conn = accepted_clients.front().lock();
accepted_clients.erase(accepted_clients.begin());
if (conn && !conn->closing)
return conn;
}
return nullptr;
}
void TCP_Server::start_accept() {
if (accept_pending || !acceptor || state != Working || !acceptor->is_open())
return;
pending_accept_socket = std::make_shared<asio::ip::tcp::socket>(io_context);
accept_pending = true;
acceptor->async_accept(
*pending_accept_socket, [this](const asio::error_code &ec) {
if (accept_pending || !acceptor || state != Working || !acceptor->is_open())
return;
pending_accept_socket = std::make_shared<asio::ip::tcp::socket>(io_context);
accept_pending = true;
acceptor->async_accept(
*pending_accept_socket, [this](const asio::error_code& ec) {
accept_pending = false;
if (state != Working)
return;
return;
if (ec == asio::error::operation_aborted)
return;
return;
if (!ec && pending_accept_socket) {
asio::error_code option_ec;
pending_accept_socket->set_option(asio::ip::tcp::no_delay(no_delay),
option_ec);
pending_accept_socket->set_option(
asio::socket_base::send_buffer_size(
static_cast<int>(connect_system_buffer_size)),
option_ec);
pending_accept_socket->set_option(
asio::socket_base::receive_buffer_size(
static_cast<int>(recv_system_buffer_size)),
option_ec);
auto conn = std::make_shared<TCP_Connect>();
conn->socket = pending_accept_socket;
conn->send_buffer.init(connect_user_buffer_size);
conn->recv_buffer.init(connect_user_buffer_size);
conn->info.fd = socket_id(conn->socket.get());
conn->info.sockaddr =
endpoint_to_sockaddr(conn->socket->remote_endpoint(option_ec));
tcp_clients[conn->info.fd] = conn;
accepted_clients.push_back(conn);
start_read(conn);
start_write(conn);
asio::error_code option_ec;
pending_accept_socket->set_option(asio::ip::tcp::no_delay(no_delay),
option_ec);
pending_accept_socket->set_option(
asio::socket_base::send_buffer_size(
static_cast<int>(connect_system_buffer_size)),
option_ec);
pending_accept_socket->set_option(
asio::socket_base::receive_buffer_size(
static_cast<int>(recv_system_buffer_size)),
option_ec);
auto conn = std::make_shared<TCP_Connect>();
conn->socket = pending_accept_socket;
conn->send_buffer.init(connect_user_buffer_size);
conn->recv_buffer.init(connect_user_buffer_size);
conn->info.fd = socket_id(conn->socket.get());
conn->info.sockaddr =
endpoint_to_sockaddr(conn->socket->remote_endpoint(option_ec));
tcp_clients[conn->info.fd] = conn;
accepted_clients.push_back(conn);
start_read(conn);
start_write(conn);
}
pending_accept_socket.reset();
start_accept();
});
});
}
void TCP_Server::start_read(const std::shared_ptr<TCP_Connect> &conn) {
if (!conn || conn->closing || conn->read_pending || !conn->socket ||
!conn->socket->is_open())
return;
conn->read_pending = true;
conn->socket->async_read_some(
asio::buffer(conn->recv_storage),
[this, conn](const asio::error_code &ec, std::size_t n) {
void TCP_Server::start_read(const std::shared_ptr<TCP_Connect>& conn) {
if (!conn || conn->closing || conn->read_pending || !conn->socket ||
!conn->socket->is_open())
return;
conn->read_pending = true;
conn->socket->async_read_some(
asio::buffer(conn->recv_storage),
[this, conn](const asio::error_code& ec, std::size_t n) {
conn->read_pending = false;
if (conn->closing)
return;
return;
if (ec == asio::error::operation_aborted)
return;
return;
if (ec == asio::error::eof || ec == asio::error::connection_reset ||
ec) {
close_client(conn);
return;
close_client(conn);
return;
}
conn->recv_buffer.write_best_effort(conn->recv_storage.data(), n);
start_read(conn);
});
});
}
void TCP_Server::start_write(const std::shared_ptr<TCP_Connect> &conn) {
if (!conn || conn->closing || conn->write_pending || !conn->socket ||
!conn->socket->is_open())
return;
conn->send_storage.resize(16 * 1024);
auto size = conn->send_buffer.peek_best_effort(conn->send_storage.data(),
conn->send_storage.size());
if (size == 0)
return;
conn->send_storage.resize(size);
conn->write_pending = true;
conn->socket->async_write_some(
asio::buffer(conn->send_storage.data(), conn->send_storage.size()),
[this, conn](const asio::error_code &ec, std::size_t sent) {
void TCP_Server::start_write(const std::shared_ptr<TCP_Connect>& conn) {
if (!conn || conn->closing || conn->write_pending || !conn->socket ||
!conn->socket->is_open())
return;
conn->send_storage.resize(16 * 1024);
auto size = conn->send_buffer.peek_best_effort(conn->send_storage.data(),
conn->send_storage.size());
if (size == 0)
return;
conn->send_storage.resize(size);
conn->write_pending = true;
conn->socket->async_write_some(
asio::buffer(conn->send_storage.data(), conn->send_storage.size()),
[this, conn](const asio::error_code& ec, std::size_t sent) {
conn->write_pending = false;
if (conn->closing)
return;
return;
if (ec == asio::error::operation_aborted)
return;
return;
if (ec) {
close_client(conn);
return;
close_client(conn);
return;
}
conn->send_buffer.skip(sent);
conn->push_speed.update(static_cast<unsigned int>(sent));
conn->send_num.update(static_cast<double>(sent));
if (sent != 0)
start_write(conn);
});
start_write(conn);
});
}
void TCP_Server::close_client(const std::shared_ptr<TCP_Connect> &conn) {
if (!conn || conn->closing)
return;
conn->closing = true;
if (conn->socket) {
asio::error_code ec;
conn->socket->cancel(ec);
conn->socket->shutdown(asio::ip::tcp::socket::shutdown_both, ec);
conn->socket->close(ec);
}
void TCP_Server::close_client(const std::shared_ptr<TCP_Connect>& conn) {
if (!conn || conn->closing)
return;
conn->closing = true;
if (conn->socket) {
asio::error_code ec;
conn->socket->cancel(ec);
conn->socket->shutdown(asio::ip::tcp::socket::shutdown_both, ec);
conn->socket->close(ec);
}
}
void TCP_Server::cleanup_closed_clients() {
for (auto it = tcp_clients.begin(); it != tcp_clients.end();) {
if (!it->second || it->second->closing || !it->second->socket ||
!it->second->socket->is_open()) {
it = tcp_clients.erase(it);
} else {
++it;
for (auto it = tcp_clients.begin(); it != tcp_clients.end();) {
if (!it->second || it->second->closing || !it->second->socket ||
!it->second->socket->is_open()) {
it = tcp_clients.erase(it);
}
else {
++it;
}
}
}
for (auto it = accepted_clients.begin(); it != accepted_clients.end();) {
auto conn = it->lock();
if (!conn || conn->closing) {
it = accepted_clients.erase(it);
} else {
++it;
for (auto it = accepted_clients.begin(); it != accepted_clients.end();) {
auto conn = it->lock();
if (!conn || conn->closing) {
it = accepted_clients.erase(it);
}
else {
++it;
}
}
}
}
void TCP_Server::flush_clients() {
start_accept();
for (auto &[_, conn] : tcp_clients) {
start_read(conn);
start_write(conn);
}
}
void TCP_Server::write_to_all_clients(std::string_view data) {
if (data.empty())
return;
for (auto &[_, conn] : tcp_clients) {
if (!conn || conn->closing)
continue;
auto written =
conn->send_buffer.write_best_effort(data.data(), data.size());
if (written < data.size())
conn->lose_speed.update(static_cast<unsigned int>(data.size() - written));
start_write(conn);
}
}
std::vector<TCP_Server::Read_Info> TCP_Server::read_from_all_clients() {
std::vector<Read_Info> ret;
std::array<char, 16 * 1024> buffer{};
for (auto &[_, conn] : tcp_clients) {
if (!conn || conn->closing)
continue;
std::string data;
for (;;) {
auto n = conn->recv_buffer.read_best_effort(buffer.data(), buffer.size());
if (n == 0)
break;
data.append(buffer.data(), n);
start_accept();
for (auto& [_, conn] : tcp_clients) {
start_read(conn);
start_write(conn);
}
if (!data.empty())
ret.push_back({conn->info, std::move(data)});
}
cleanup_closed_clients();
return ret;
}
void TCP_Server::write_to_all_clients(std::string_view data) {
if (data.empty())
return;
for (auto& [_, conn] : tcp_clients) {
if (!conn || conn->closing)
continue;
auto written =
conn->send_buffer.write_best_effort(data.data(), data.size());
if (written < data.size())
conn->lose_speed.update(static_cast<unsigned int>(data.size() - written));
start_write(conn);
}
}
std::vector<TCP_Server::Read_Info> TCP_Server::read_from_all_clients() {
std::vector<Read_Info> ret;
std::array<char, 16 * 1024> buffer{};
for (auto& [_, conn] : tcp_clients) {
if (!conn || conn->closing)
continue;
std::string data;
for (;;) {
auto n = conn->recv_buffer.read_best_effort(buffer.data(), buffer.size());
if (n == 0)
break;
data.append(buffer.data(), n);
}
if (!data.empty())
ret.push_back({conn->info, std::move(data)});
}
cleanup_closed_clients();
return ret;
}
std::vector<Socket_FD> TCP_Server::client_fds() {
std::vector<Socket_FD> ret;
for (const auto &[fd, conn] : tcp_clients) {
if (conn && !conn->closing)
ret.push_back(fd);
}
return ret;
std::vector<Socket_FD> ret;
for (const auto& [fd, conn] : tcp_clients) {
if (conn && !conn->closing)
ret.push_back(fd);
}
return ret;
}
std::vector<std::shared_ptr<TCP_Connect>> TCP_Server::get_all_clients() {
std::vector<std::shared_ptr<TCP_Connect>> ret;
for (const auto &[_, conn] : tcp_clients) {
if (conn && !conn->closing)
ret.push_back(conn);
}
return ret;
std::vector<std::shared_ptr<TCP_Connect>> ret;
for (const auto& [_, conn] : tcp_clients) {
if (conn && !conn->closing)
ret.push_back(conn);
}
return ret;
}
std::string TCP_Server::to_string() {
return "TCP_Server:[" + (address ? address->to_string() : std::string{}) +
"]";
return "TCP_Server:[" + (address ? address->to_string() : std::string{}) +
"]";
}
} // namespace Psc::asio_socket
+61 -68
View File
@@ -1,82 +1,75 @@
#pragma once
#include <string_view>
#include "ASIO_Utils.h"
#include "Core/Statistics/Statistics.h"
#include <array>
#include <map>
#include <memory>
#include <string_view>
#include <vector>
#include "ASIO_Utils.h"
#include "Core/Statistics/Statistics.h"
namespace Psc::asio_socket {
class TCP_Connect {
public:
Accept_Info info;
std::shared_ptr<asio::ip::tcp::socket> socket;
ASIO_StreamRingBuffer send_buffer;
ASIO_StreamRingBuffer recv_buffer;
Speed_Statistics push_speed;
Speed_Statistics lose_speed;
Value_Statistics send_num;
std::array<char, 16 * 1024> recv_storage{};
std::vector<char> send_storage;
bool read_pending = false;
bool write_pending = false;
bool closing = false;
[[nodiscard]] std::string to_string() const { return info.to_string(); }
Accept_Info info;
std::shared_ptr<asio::ip::tcp::socket> socket;
ASIO_StreamRingBuffer send_buffer;
ASIO_StreamRingBuffer recv_buffer;
Speed_Statistics push_speed;
Speed_Statistics lose_speed;
Value_Statistics send_num;
std::array<char, 16 * 1024> recv_storage{};
std::vector<char> send_storage;
bool read_pending = false;
bool write_pending = false;
bool closing = false;
[[nodiscard]] std::string to_string() const {
return info.to_string();
}
};
class TCP_Server : public Socket_Base {
public:
~TCP_Server() override;
void tick();
TCP_Server &set_tcp_no_delay(bool no_delay);
void close() override;
TCP_Server &create();
TCP_Server &set_connect_system_buffer_size(size_t size);
TCP_Server &set_connect_user_buffer_size(size_t size);
TCP_Server &set_recv_system_buffer_size(size_t size);
bool listen(std::string_view ip, std::uint32_t port);
bool listen(const Sockaddr_In &address);
std::string to_string() override;
void write_to_all_clients(std::string_view data);
struct Read_Info {
Accept_Info info;
std::string data;
};
enum State { Not_Created, Not_Set_Field, Not_bind, Not_Listen, Working };
State state = Not_Created;
std::vector<Read_Info> read_from_all_clients();
std::vector<Socket_FD> client_fds();
std::vector<std::shared_ptr<TCP_Connect>> get_all_clients();
void flush_clients();
[[nodiscard]] std::shared_ptr<TCP_Connect> accept() const;
~TCP_Server() override;
void tick();
TCP_Server& set_tcp_no_delay(bool no_delay);
void close() override;
TCP_Server& create();
TCP_Server& set_connect_system_buffer_size(size_t size);
TCP_Server& set_connect_user_buffer_size(size_t size);
TCP_Server& set_recv_system_buffer_size(size_t size);
bool listen(std::string_view ip, std::uint32_t port);
bool listen(const Sockaddr_In& address);
std::string to_string() override;
void write_to_all_clients(std::string_view data);
struct Read_Info {
Accept_Info info;
std::string data;
};
enum State { Not_Created,
Not_Set_Field,
Not_bind,
Not_Listen,
Working };
State state = Not_Created;
std::vector<Read_Info> read_from_all_clients();
std::vector<Socket_FD> client_fds();
std::vector<std::shared_ptr<TCP_Connect>> get_all_clients();
void flush_clients();
[[nodiscard]] std::shared_ptr<TCP_Connect> accept() const;
protected:
mutable asio::io_context io_context;
std::unique_ptr<asio::ip::tcp::acceptor> acceptor;
std::map<Socket_FD, std::shared_ptr<TCP_Connect>> tcp_clients;
size_t connect_user_buffer_size = 1024 * 1024;
size_t connect_system_buffer_size = 4096 * 10;
size_t max_pending_send_bytes = 1024 * 1024;
size_t recv_system_buffer_size = 1024 * 1024;
bool no_delay = true;
bool accept_pending = false;
std::shared_ptr<asio::ip::tcp::socket> pending_accept_socket;
mutable std::vector<std::weak_ptr<TCP_Connect>> accepted_clients;
void start_accept();
void start_read(const std::shared_ptr<TCP_Connect> &conn);
void start_write(const std::shared_ptr<TCP_Connect> &conn);
void close_client(const std::shared_ptr<TCP_Connect> &conn);
void cleanup_closed_clients();
mutable asio::io_context io_context;
std::unique_ptr<asio::ip::tcp::acceptor> acceptor;
std::map<Socket_FD, std::shared_ptr<TCP_Connect>> tcp_clients;
size_t connect_user_buffer_size = 1024 * 1024;
size_t connect_system_buffer_size = 4096 * 10;
size_t max_pending_send_bytes = 1024 * 1024;
size_t recv_system_buffer_size = 1024 * 1024;
bool no_delay = true;
bool accept_pending = false;
std::shared_ptr<asio::ip::tcp::socket> pending_accept_socket;
mutable std::vector<std::weak_ptr<TCP_Connect>> accepted_clients;
void start_accept();
void start_read(const std::shared_ptr<TCP_Connect>& conn);
void start_write(const std::shared_ptr<TCP_Connect>& conn);
void close_client(const std::shared_ptr<TCP_Connect>& conn);
void cleanup_closed_clients();
};
} // namespace Psc::asio_socket
+143 -164
View File
@@ -1,173 +1,152 @@
#include "UDP_Client.h"
#include <string_view>
namespace Psc::asio_socket {
void UDP_Client::create() {
socket = std::make_unique<asio::ip::udp::socket>(io_context);
asio::error_code ec;
socket->open(asio::ip::udp::v4(), ec);
connected = false;
connect_pending = false;
read_pending = false;
write_pending = false;
recv_storage.assign(static_cast<size_t>(read_chunk_size), 0);
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
send_buffer.init(user_buffer_size);
recv_buffer.init(user_buffer_size);
state = Not_Set_Field;
}
bool UDP_Client::connect() {
if (connected || connect_pending)
return true;
if (!socket)
create();
if (dest_address.ip.empty() || dest_address.port == 0) {
state = Not_Set_Field;
return false;
}
asio::error_code ec;
endpoint =
asio::ip::udp::endpoint(asio::ip::make_address(dest_address.ip, ec),
static_cast<unsigned short>(dest_address.port));
if (ec) {
state = Not_Set_Field;
return false;
}
io_context.restart();
connect_pending = true;
state = Connecting;
socket->async_connect(endpoint, [this](const asio::error_code &connect_ec) {
connect_pending = false;
if (connect_ec == asio::error::operation_aborted)
return;
if (connect_ec) {
connected = false;
state = Not_Set_Field;
return;
}
connected = true;
state = Working;
start_read();
start_write();
});
return true;
}
std::string UDP_Client::read() {
std::string ret(static_cast<size_t>(read_chunk_size), '\0');
std::size_t out_len = ret.size();
if (!recv_buffer.read(ret.data(), out_len))
return {};
ret.resize(out_len);
return ret;
}
void UDP_Client::tick() {
if (!connected && !connect_pending && !dest_address.ip.empty() &&
dest_address.port != 0) {
connect();
}
if (connected) {
start_read();
start_write();
}
io_context.restart();
io_context.poll();
if (connected) {
start_read();
start_write();
}
}
void UDP_Client::send(std::string_view data) {
if (data.empty())
return;
send_buffer.write(data.data(), data.size());
if (!connected)
connect();
if (connected)
start_write();
}
void UDP_Client::start_read() {
if (read_pending || !connected || !socket || !socket->is_open())
return;
if (recv_storage.empty())
recv_storage.assign(static_cast<size_t>(read_chunk_size), 0);
read_pending = true;
socket->async_receive(asio::buffer(recv_storage.data(), recv_storage.size()),
[this](const asio::error_code &ec, std::size_t n) {
read_pending = false;
if (!connected)
return;
if (ec == asio::error::operation_aborted)
return;
if (ec) {
connected = false;
state = Not_Set_Field;
return;
}
recv_buffer.write(recv_storage.data(), n);
start_read();
});
}
void UDP_Client::start_write() {
if (write_pending || !connected || !socket || !socket->is_open())
return;
if (send_storage.empty())
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
std::size_t out_len = send_storage.size();
if (!send_buffer.peek(send_storage.data(), out_len))
return;
send_storage.resize(out_len);
write_pending = true;
socket->async_send(asio::buffer(send_storage.data(), send_storage.size()),
[this](const asio::error_code &ec, std::size_t) {
write_pending = false;
if (!connected)
return;
if (ec == asio::error::operation_aborted)
return;
if (ec) {
connected = false;
state = Not_Set_Field;
return;
}
send_buffer.skip_one();
send_storage.assign(static_cast<size_t>(read_chunk_size),
0);
start_write();
});
}
void UDP_Client::close() {
connected = false;
connect_pending = false;
read_pending = false;
write_pending = false;
if (socket) {
socket = std::make_unique<asio::ip::udp::socket>(io_context);
asio::error_code ec;
socket->cancel(ec);
socket->close(ec);
}
state = Not_Created;
socket->open(asio::ip::udp::v4(), ec);
connected = false;
connect_pending = false;
read_pending = false;
write_pending = false;
recv_storage.assign(static_cast<size_t>(read_chunk_size), 0);
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
send_buffer.init(user_buffer_size);
recv_buffer.init(user_buffer_size);
state = Not_Set_Field;
}
bool UDP_Client::connect() {
if (connected || connect_pending)
return true;
if (!socket)
create();
if (dest_address.ip.empty() || dest_address.port == 0) {
state = Not_Set_Field;
return false;
}
asio::error_code ec;
endpoint =
asio::ip::udp::endpoint(asio::ip::make_address(dest_address.ip, ec),
static_cast<unsigned short>(dest_address.port));
if (ec) {
state = Not_Set_Field;
return false;
}
io_context.restart();
connect_pending = true;
state = Connecting;
socket->async_connect(endpoint, [this](const asio::error_code& connect_ec) {
connect_pending = false;
if (connect_ec == asio::error::operation_aborted)
return;
if (connect_ec) {
connected = false;
state = Not_Set_Field;
return;
}
connected = true;
state = Working;
start_read();
start_write();
});
return true;
}
std::string UDP_Client::read() {
std::string ret(static_cast<size_t>(read_chunk_size), '\0');
std::size_t out_len = ret.size();
if (!recv_buffer.read(ret.data(), out_len))
return {};
ret.resize(out_len);
return ret;
}
void UDP_Client::tick() {
if (!connected && !connect_pending && !dest_address.ip.empty() &&
dest_address.port != 0) {
connect();
}
if (connected) {
start_read();
start_write();
}
io_context.restart();
io_context.poll();
if (connected) {
start_read();
start_write();
}
}
void UDP_Client::send(std::string_view data) {
if (data.empty())
return;
send_buffer.write(data.data(), data.size());
if (!connected)
connect();
if (connected)
start_write();
}
void UDP_Client::start_read() {
if (read_pending || !connected || !socket || !socket->is_open())
return;
if (recv_storage.empty())
recv_storage.assign(static_cast<size_t>(read_chunk_size), 0);
read_pending = true;
socket->async_receive(asio::buffer(recv_storage.data(), recv_storage.size()),
[this](const asio::error_code& ec, std::size_t n) {
read_pending = false;
if (!connected)
return;
if (ec == asio::error::operation_aborted)
return;
if (ec) {
connected = false;
state = Not_Set_Field;
return;
}
recv_buffer.write(recv_storage.data(), n);
start_read();
});
}
void UDP_Client::start_write() {
if (write_pending || !connected || !socket || !socket->is_open())
return;
if (send_storage.empty())
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
std::size_t out_len = send_storage.size();
if (!send_buffer.peek(send_storage.data(), out_len))
return;
send_storage.resize(out_len);
write_pending = true;
socket->async_send(asio::buffer(send_storage.data(), send_storage.size()),
[this](const asio::error_code& ec, std::size_t) {
write_pending = false;
if (!connected)
return;
if (ec == asio::error::operation_aborted)
return;
if (ec) {
connected = false;
state = Not_Set_Field;
return;
}
send_buffer.skip_one();
send_storage.assign(static_cast<size_t>(read_chunk_size),
0);
start_write();
});
}
void UDP_Client::close() {
connected = false;
connect_pending = false;
read_pending = false;
write_pending = false;
if (socket) {
asio::error_code ec;
socket->cancel(ec);
socket->close(ec);
}
state = Not_Created;
}
std::string UDP_Client::to_string() {
return "UDP_Client:[" + dest_address.to_string() + "]";
return "UDP_Client:[" + dest_address.to_string() + "]";
}
} // namespace Psc::asio_socket
+38 -38
View File
@@ -1,48 +1,48 @@
#pragma once
#include <string_view>
#include "ASIO_Utils.h"
#include <array>
#include <memory>
#include <string_view>
#include <vector>
#include "ASIO_Utils.h"
namespace Psc::asio_socket {
class UDP_Client {
public:
void set_dest_address(const Sockaddr_In& value) {
dest_address = value;
}
void set_dest_address(std::string_view ip, std::uint32_t port) {
dest_address = Sockaddr_In(ip, port);
}
void create();
bool connect();
std::string read();
void tick();
void send(std::string_view data);
void close();
std::string to_string();
Sockaddr_In dest_address;
enum State {
Not_Created,
Not_Set_Field,
Connecting,
Working,
};
State state = Not_Created;
void set_dest_address(const Sockaddr_In& value) {
dest_address = value;
}
void set_dest_address(std::string_view ip, std::uint32_t port) {
dest_address = Sockaddr_In(ip, port);
}
void create();
bool connect();
std::string read();
void tick();
void send(std::string_view data);
void close();
std::string to_string();
Sockaddr_In dest_address;
enum State {
Not_Created,
Not_Set_Field,
Connecting,
Working,
};
State state = Not_Created;
private:
asio::io_context io_context;
std::unique_ptr<asio::ip::udp::socket> socket;
asio::ip::udp::endpoint endpoint;
bool connected = false;
int read_chunk_size = 1024 * 1024;
size_t user_buffer_size = 1024 * 1024;
ASIO_PacketRingBuffer send_buffer;
ASIO_PacketRingBuffer recv_buffer;
std::vector<char> recv_storage;
std::vector<char> send_storage;
bool connect_pending = false;
bool read_pending = false;
bool write_pending = false;
void start_read();
void start_write();
asio::io_context io_context;
std::unique_ptr<asio::ip::udp::socket> socket;
asio::ip::udp::endpoint endpoint;
bool connected = false;
int read_chunk_size = 1024 * 1024;
size_t user_buffer_size = 1024 * 1024;
ASIO_PacketRingBuffer send_buffer;
ASIO_PacketRingBuffer recv_buffer;
std::vector<char> recv_storage;
std::vector<char> send_storage;
bool connect_pending = false;
bool read_pending = false;
bool write_pending = false;
void start_read();
void start_write();
};
} // namespace Psc::asio_socket
+129 -153
View File
@@ -1,193 +1,169 @@
#include "UDP_Server.h"
#include <algorithm>
#include <string_view>
namespace Psc::asio_socket {
void UDP_Server::create() {
socket = std::make_unique<asio::ip::udp::socket>(io_context);
socket_fd = socket_id(socket.get());
recv_buffer.init(buffer_size);
send_buffer.init(buffer_size);
recv_storage.assign(static_cast<size_t>(read_chunk_size), 0);
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
read_pending = false;
write_pending = false;
state = Not_Bind;
}
bool UDP_Server::bind() {
if (!socket)
create();
if (bind_address.ip.empty() || bind_address.port == 0) {
state = Not_Set_Field;
return false;
}
asio::error_code ec;
auto address_value = asio::ip::make_address(bind_address.ip, ec);
if (ec) {
state = Not_Bind;
return false;
}
asio::ip::udp::endpoint endpoint(
address_value, static_cast<unsigned short>(bind_address.port));
socket->open(endpoint.protocol(), ec);
if (ec)
return false;
socket->set_option(asio::socket_base::reuse_address(true), ec);
socket->set_option(
asio::socket_base::receive_buffer_size(static_cast<int>(buffer_size)),
ec);
socket->bind(endpoint, ec);
if (ec) {
state = Not_Bind;
return false;
}
bound = true;
state = Working;
start_read();
return true;
}
void UDP_Server::tick() {
if (!bound || !socket)
return;
start_read();
start_write();
io_context.restart();
io_context.poll();
start_read();
start_write();
}
void UDP_Server::start_read() {
if (read_pending || !bound || !socket || !socket->is_open())
return;
if (recv_storage.empty())
socket = std::make_unique<asio::ip::udp::socket>(io_context);
socket_fd = socket_id(socket.get());
recv_buffer.init(buffer_size);
send_buffer.init(buffer_size);
recv_storage.assign(static_cast<size_t>(read_chunk_size), 0);
read_pending = true;
socket->async_receive_from(
asio::buffer(recv_storage.data(), recv_storage.size()), recv_endpoint,
[this](const asio::error_code &ec, std::size_t n) {
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
read_pending = false;
write_pending = false;
state = Not_Bind;
}
bool UDP_Server::bind() {
if (!socket)
create();
if (bind_address.ip.empty() || bind_address.port == 0) {
state = Not_Set_Field;
return false;
}
asio::error_code ec;
auto address_value = asio::ip::make_address(bind_address.ip, ec);
if (ec) {
state = Not_Bind;
return false;
}
asio::ip::udp::endpoint endpoint(
address_value, static_cast<unsigned short>(bind_address.port));
socket->open(endpoint.protocol(), ec);
if (ec)
return false;
socket->set_option(asio::socket_base::reuse_address(true), ec);
socket->set_option(
asio::socket_base::receive_buffer_size(static_cast<int>(buffer_size)),
ec);
socket->bind(endpoint, ec);
if (ec) {
state = Not_Bind;
return false;
}
bound = true;
state = Working;
start_read();
return true;
}
void UDP_Server::tick() {
if (!bound || !socket)
return;
start_read();
start_write();
io_context.restart();
io_context.poll();
start_read();
start_write();
}
void UDP_Server::start_read() {
if (read_pending || !bound || !socket || !socket->is_open())
return;
if (recv_storage.empty())
recv_storage.assign(static_cast<size_t>(read_chunk_size), 0);
read_pending = true;
socket->async_receive_from(
asio::buffer(recv_storage.data(), recv_storage.size()), recv_endpoint,
[this](const asio::error_code& ec, std::size_t n) {
read_pending = false;
if (!bound)
return;
return;
if (ec == asio::error::operation_aborted)
return;
return;
if (ec)
return;
return;
last_peer = endpoint_to_sockaddr(recv_endpoint);
has_last_peer = true;
if (std::find(clients.begin(), clients.end(), last_peer) ==
clients.end()) {
clients.push_back(last_peer);
clients.push_back(last_peer);
}
if (recv_buffer.write(recv_storage.data(), n)) {
recv_peers.push_back(last_peer);
recv_peers.push_back(last_peer);
}
start_read();
});
});
}
void UDP_Server::start_write() {
if (write_pending || !bound || !socket || !socket->is_open())
return;
if (send_endpoints.empty())
return;
if (send_storage.empty())
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
std::size_t out_len = send_storage.size();
if (!send_buffer.peek(send_storage.data(), out_len))
return;
send_storage.resize(out_len);
auto endpoint = send_endpoints.front();
write_pending = true;
socket->async_send_to(
asio::buffer(send_storage.data(), send_storage.size()), endpoint,
[this](const asio::error_code &ec, std::size_t) {
if (write_pending || !bound || !socket || !socket->is_open())
return;
if (send_endpoints.empty())
return;
if (send_storage.empty())
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
std::size_t out_len = send_storage.size();
if (!send_buffer.peek(send_storage.data(), out_len))
return;
send_storage.resize(out_len);
auto endpoint = send_endpoints.front();
write_pending = true;
socket->async_send_to(
asio::buffer(send_storage.data(), send_storage.size()), endpoint,
[this](const asio::error_code& ec, std::size_t) {
write_pending = false;
if (!bound)
return;
return;
if (ec == asio::error::operation_aborted)
return;
return;
if (!ec) {
send_buffer.skip_one();
if (!send_endpoints.empty())
send_endpoints.erase(send_endpoints.begin());
send_buffer.skip_one();
if (!send_endpoints.empty())
send_endpoints.erase(send_endpoints.begin());
}
send_storage.assign(static_cast<size_t>(read_chunk_size), 0);
start_write();
});
});
}
void UDP_Server::close() {
read_pending = false;
write_pending = false;
if (socket) {
asio::error_code ec;
socket->cancel(ec);
socket->close(ec);
}
clients.clear();
recv_peers.clear();
send_endpoints.clear();
bound = false;
has_last_peer = false;
socket_fd = static_cast<Socket_FD>(-1);
state = Not_Created;
read_pending = false;
write_pending = false;
if (socket) {
asio::error_code ec;
socket->cancel(ec);
socket->close(ec);
}
clients.clear();
recv_peers.clear();
send_endpoints.clear();
bound = false;
has_last_peer = false;
socket_fd = static_cast<Socket_FD>(-1);
state = Not_Created;
}
std::vector<UDP_Server::Read_Info> UDP_Server::read() {
std::vector<Read_Info> ret;
std::string data(static_cast<size_t>(read_chunk_size), '\0');
while (!recv_peers.empty()) {
std::size_t out_len = data.size();
if (!recv_buffer.read(data.data(), out_len))
break;
ret.push_back({recv_peers.front(), std::string(data.data(), out_len)});
recv_peers.erase(recv_peers.begin());
}
return ret;
std::vector<Read_Info> ret;
std::string data(static_cast<size_t>(read_chunk_size), '\0');
while (!recv_peers.empty()) {
std::size_t out_len = data.size();
if (!recv_buffer.read(data.data(), out_len))
break;
ret.push_back({recv_peers.front(), std::string(data.data(), out_len)});
recv_peers.erase(recv_peers.begin());
}
return ret;
}
void UDP_Server::reply_last_peer(std::string_view data) {
if (!has_last_peer)
return;
send_to(last_peer.ip, static_cast<uint16_t>(last_peer.port), data);
if (!has_last_peer)
return;
send_to(last_peer.ip, static_cast<uint16_t>(last_peer.port), data);
}
void UDP_Server::send_to(std::string_view ip, uint16_t port,
std::string_view data) {
if (data.empty())
return;
if (!bound && !bind())
return;
asio::error_code ec;
auto endpoint = asio::ip::udp::endpoint(asio::ip::make_address(ip, ec), port);
if (ec)
return;
if (send_buffer.write(data.data(), data.size())) {
send_endpoints.push_back(endpoint);
}
start_write();
if (data.empty())
return;
if (!bound && !bind())
return;
asio::error_code ec;
auto endpoint = asio::ip::udp::endpoint(asio::ip::make_address(ip, ec), port);
if (ec)
return;
if (send_buffer.write(data.data(), data.size())) {
send_endpoints.push_back(endpoint);
}
start_write();
}
void UDP_Server::write_to_all_clients(std::string_view msg) {
for (const auto &client : clients) {
send_to(client.ip, static_cast<uint16_t>(client.port), msg);
}
for (const auto& client : clients) {
send_to(client.ip, static_cast<uint16_t>(client.port), msg);
}
}
} // namespace Psc::asio_socket
+51 -61
View File
@@ -1,71 +1,61 @@
#pragma once
#include <string_view>
#include "ASIO_Utils.h"
#include <deque>
#include <memory>
#include <string_view>
#include <vector>
#include "ASIO_Utils.h"
namespace Psc::asio_socket {
class UDP_Server {
public:
void set_bind_address(const Sockaddr_In &addr) { bind_address = addr; }
void set_bind_address(std::string_view ip, std::uint32_t port) {
bind_address = Sockaddr_In(ip, port);
}
[[nodiscard]] std::string to_string() const {
return bind_address.to_string();
}
void create();
bool bind();
void tick();
void close();
void reply_last_peer(std::string_view data);
void send_to(std::string_view ip, uint16_t port, std::string_view data);
void write_to_all_clients(std::string_view msg);
struct Read_Info {
Sockaddr_In peer;
std::string data;
};
std::vector<Read_Info> read();
enum State {
Not_Created,
Not_Set_Field,
Not_Bind,
Working,
};
State state = Not_Created;
Sockaddr_In bind_address;
std::vector<Sockaddr_In> clients;
Socket_FD socket_fd = static_cast<Socket_FD>(-1);
bool bound = false;
int read_chunk_size = 1024 * 1024;
Sockaddr_In last_peer{};
bool has_last_peer = false;
size_t buffer_size = 4096 * 10;
void set_bind_address(const Sockaddr_In& addr) {
bind_address = addr;
}
void set_bind_address(std::string_view ip, std::uint32_t port) {
bind_address = Sockaddr_In(ip, port);
}
[[nodiscard]] std::string to_string() const {
return bind_address.to_string();
}
void create();
bool bind();
void tick();
void close();
void reply_last_peer(std::string_view data);
void send_to(std::string_view ip, uint16_t port, std::string_view data);
void write_to_all_clients(std::string_view msg);
struct Read_Info {
Sockaddr_In peer;
std::string data;
};
std::vector<Read_Info> read();
enum State {
Not_Created,
Not_Set_Field,
Not_Bind,
Working,
};
State state = Not_Created;
Sockaddr_In bind_address;
std::vector<Sockaddr_In> clients;
Socket_FD socket_fd = static_cast<Socket_FD>(-1);
bool bound = false;
int read_chunk_size = 1024 * 1024;
Sockaddr_In last_peer{};
bool has_last_peer = false;
size_t buffer_size = 4096 * 10;
private:
asio::io_context io_context;
std::unique_ptr<asio::ip::udp::socket> socket;
asio::ip::udp::endpoint recv_endpoint;
ASIO_PacketRingBuffer recv_buffer;
ASIO_PacketRingBuffer send_buffer;
std::vector<Sockaddr_In> recv_peers;
std::vector<asio::ip::udp::endpoint> send_endpoints;
std::vector<char> recv_storage;
std::vector<char> send_storage;
bool read_pending = false;
bool write_pending = false;
void start_read();
void start_write();
asio::io_context io_context;
std::unique_ptr<asio::ip::udp::socket> socket;
asio::ip::udp::endpoint recv_endpoint;
ASIO_PacketRingBuffer recv_buffer;
ASIO_PacketRingBuffer send_buffer;
std::vector<Sockaddr_In> recv_peers;
std::vector<asio::ip::udp::endpoint> send_endpoints;
std::vector<char> recv_storage;
std::vector<char> send_storage;
bool read_pending = false;
bool write_pending = false;
void start_read();
void start_write();
};
} // namespace Psc::asio_socket
-1
View File
@@ -1,5 +1,4 @@
#pragma once
#include "Socket.h"
#include "Socket_Coro.h"
#include "TCP_Client.h"
+1 -4
View File
@@ -1,9 +1,6 @@
#pragma once
#if defined(_WIN32) && !defined(_WIN32_WINNT)
#define _WIN32_WINNT 0x0601
#endif
#include "ByteOrder.h"
#include <asio.hpp>
#include "ByteOrder.h"