concurrencpp 换成asio协程
This commit is contained in:
+61
-23
@@ -1,6 +1,11 @@
|
|||||||
#pragma once
|
#pragma once
|
||||||
#include <concurrencpp/concurrencpp.h>
|
#include <asio/as_tuple.hpp>
|
||||||
|
#include <asio/awaitable.hpp>
|
||||||
|
#include <asio/dispatch.hpp>
|
||||||
|
#include <asio/this_coro.hpp>
|
||||||
|
#include <asio/use_awaitable.hpp>
|
||||||
#include <exception>
|
#include <exception>
|
||||||
|
#include <functional>
|
||||||
#include <memory>
|
#include <memory>
|
||||||
#include <type_traits>
|
#include <type_traits>
|
||||||
#include <utility>
|
#include <utility>
|
||||||
@@ -8,39 +13,72 @@
|
|||||||
namespace Psc::coro {
|
namespace Psc::coro {
|
||||||
template <typename T> class Result_Completion {
|
template <typename T> class Result_Completion {
|
||||||
public:
|
public:
|
||||||
explicit Result_Completion(
|
using Complete = std::function<void(std::exception_ptr, T)>;
|
||||||
std::shared_ptr<concurrencpp::result_promise<T>> promise)
|
explicit Result_Completion(Complete complete) : complete_(std::move(complete)) {}
|
||||||
: promise_(std::move(promise)) {}
|
void operator()(T value) { complete_(nullptr, std::move(value)); }
|
||||||
void operator()(T value) { promise_->set_result(std::move(value)); }
|
|
||||||
void set_exception(std::exception_ptr exception) {
|
void set_exception(std::exception_ptr exception) {
|
||||||
promise_->set_exception(exception);
|
complete_(exception, T{});
|
||||||
}
|
}
|
||||||
|
|
||||||
private:
|
private:
|
||||||
std::shared_ptr<concurrencpp::result_promise<T>> promise_;
|
Complete complete_;
|
||||||
};
|
};
|
||||||
template <> class Result_Completion<void> {
|
template <> class Result_Completion<void> {
|
||||||
public:
|
public:
|
||||||
explicit Result_Completion(
|
using Complete = std::function<void(std::exception_ptr)>;
|
||||||
std::shared_ptr<concurrencpp::result_promise<void>> promise)
|
explicit Result_Completion(Complete complete) : complete_(std::move(complete)) {}
|
||||||
: promise_(std::move(promise)) {}
|
void operator()() { complete_(nullptr); }
|
||||||
void operator()() { promise_->set_result(); }
|
|
||||||
void set_exception(std::exception_ptr exception) {
|
void set_exception(std::exception_ptr exception) {
|
||||||
promise_->set_exception(exception);
|
complete_(exception);
|
||||||
}
|
}
|
||||||
|
|
||||||
private:
|
private:
|
||||||
std::shared_ptr<concurrencpp::result_promise<void>> promise_;
|
Complete complete_;
|
||||||
};
|
};
|
||||||
template <typename T, typename Starter>
|
template <typename T, typename Starter>
|
||||||
[[nodiscard]] concurrencpp::result<T> callback_result(Starter &&starter) {
|
[[nodiscard]] asio::awaitable<T> callback_result(Starter &&starter) {
|
||||||
auto promise = std::make_shared<concurrencpp::result_promise<T>>();
|
auto executor = co_await asio::this_coro::executor;
|
||||||
auto result = promise->get_result();
|
auto token = asio::as_tuple(asio::use_awaitable);
|
||||||
try {
|
auto [exception, value] = co_await asio::async_initiate<decltype(token), void(std::exception_ptr, T)>(
|
||||||
std::forward<Starter>(starter)(Result_Completion<T>{promise});
|
[starter = std::forward<Starter>(starter), executor](auto handler) mutable {
|
||||||
} catch (...) {
|
auto handler_ptr = std::make_shared<std::decay_t<decltype(handler)>>(std::move(handler));
|
||||||
promise->set_exception(std::current_exception());
|
auto complete = [handler_ptr, executor](std::exception_ptr exception, T value) mutable {
|
||||||
|
asio::dispatch(executor, [handler_ptr, exception, value = std::move(value)]() mutable {
|
||||||
|
(*handler_ptr)(exception, std::move(value));
|
||||||
|
});
|
||||||
|
};
|
||||||
|
try {
|
||||||
|
starter(Result_Completion<T>{std::move(complete)});
|
||||||
|
} catch (...) {
|
||||||
|
complete(std::current_exception(), T{});
|
||||||
|
}
|
||||||
|
},
|
||||||
|
token);
|
||||||
|
if (exception) {
|
||||||
|
std::rethrow_exception(exception);
|
||||||
}
|
}
|
||||||
return result;
|
co_return std::move(value);
|
||||||
|
}
|
||||||
|
template <typename Starter>
|
||||||
|
[[nodiscard]] asio::awaitable<void> callback_result(Starter &&starter) {
|
||||||
|
auto executor = co_await asio::this_coro::executor;
|
||||||
|
auto token = asio::as_tuple(asio::use_awaitable);
|
||||||
|
auto [exception] = co_await asio::async_initiate<decltype(token), void(std::exception_ptr)>(
|
||||||
|
[starter = std::forward<Starter>(starter), executor](auto handler) mutable {
|
||||||
|
auto handler_ptr = std::make_shared<std::decay_t<decltype(handler)>>(std::move(handler));
|
||||||
|
auto complete = [handler_ptr, executor](std::exception_ptr exception) mutable {
|
||||||
|
asio::dispatch(executor, [handler_ptr, exception]() mutable {
|
||||||
|
(*handler_ptr)(exception);
|
||||||
|
});
|
||||||
|
};
|
||||||
|
try {
|
||||||
|
starter(Result_Completion<void>{std::move(complete)});
|
||||||
|
} catch (...) {
|
||||||
|
complete(std::current_exception());
|
||||||
|
}
|
||||||
|
},
|
||||||
|
token);
|
||||||
|
if (exception) {
|
||||||
|
std::rethrow_exception(exception);
|
||||||
|
}
|
||||||
|
co_return;
|
||||||
}
|
}
|
||||||
} // namespace Psc::coro
|
} // namespace Psc::coro
|
||||||
|
|||||||
+39
-34
@@ -1,11 +1,11 @@
|
|||||||
#pragma once
|
#pragma once
|
||||||
|
|
||||||
#include "Serial.h"
|
#include "Serial.h"
|
||||||
|
#include <asio/awaitable.hpp>
|
||||||
#include <asio/error_code.hpp>
|
#include <asio/error_code.hpp>
|
||||||
#include <asio/io_context.hpp>
|
#include <asio/io_context.hpp>
|
||||||
#include <asio/serial_port.hpp>
|
#include <asio/serial_port.hpp>
|
||||||
#include <asio/write.hpp>
|
#include <asio/write.hpp>
|
||||||
#include <concurrencpp/concurrencpp.h>
|
|
||||||
|
|
||||||
#include <cstddef>
|
#include <cstddef>
|
||||||
#include <memory>
|
#include <memory>
|
||||||
@@ -50,53 +50,28 @@ public:
|
|||||||
io_context_.restart();
|
io_context_.restart();
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> tick_coro() {
|
[[nodiscard]] asio::awaitable<void> tick_coro() {
|
||||||
io_context_.poll();
|
tick();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<std::string>
|
[[nodiscard]] asio::awaitable<std::string>
|
||||||
read_coro(std::size_t max_size = 16 * 1024) {
|
read_coro(std::size_t max_size = 16 * 1024) {
|
||||||
co_await tick_coro();
|
co_return read_impl(max_size);
|
||||||
if (!port_ || !port_->is_open()) {
|
|
||||||
co_return "";
|
|
||||||
}
|
|
||||||
|
|
||||||
auto available = serial::get_available_bytes(port_->native_handle());
|
|
||||||
if (available <= 0) {
|
|
||||||
co_return "";
|
|
||||||
}
|
|
||||||
|
|
||||||
auto data = serial::read_all(port_->native_handle());
|
|
||||||
if (data.size() > max_size) {
|
|
||||||
data.resize(max_size);
|
|
||||||
}
|
|
||||||
co_return data;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<std::size_t> write_coro(std::string data) {
|
[[nodiscard]] asio::awaitable<std::size_t> write_coro(std::string data) {
|
||||||
co_await tick_coro();
|
co_return write_impl(data);
|
||||||
if (!port_ || !port_->is_open() || data.empty()) {
|
|
||||||
co_return 0;
|
|
||||||
}
|
|
||||||
|
|
||||||
asio::error_code ec;
|
|
||||||
auto size = asio::write(*port_, asio::buffer(data), ec);
|
|
||||||
if (ec) {
|
|
||||||
co_return 0;
|
|
||||||
}
|
|
||||||
|
|
||||||
co_return size;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
std::string read(int64_t size = -1) {
|
std::string read(int64_t size = -1) {
|
||||||
auto max_size = size > 0 ? static_cast<std::size_t>(size)
|
auto max_size = size > 0 ? static_cast<std::size_t>(size)
|
||||||
: static_cast<std::size_t>(16 * 1024);
|
: static_cast<std::size_t>(16 * 1024);
|
||||||
return (read_coro(max_size)).get();
|
return read_impl(max_size);
|
||||||
}
|
}
|
||||||
|
|
||||||
int64_t write(const std::string &data) {
|
int64_t write(const std::string &data) {
|
||||||
return static_cast<int64_t>(write_coro(data).get());
|
return static_cast<int64_t>(write_impl(data));
|
||||||
}
|
}
|
||||||
|
|
||||||
int get_available_bytes() {
|
int get_available_bytes() {
|
||||||
@@ -107,6 +82,36 @@ public:
|
|||||||
}
|
}
|
||||||
|
|
||||||
private:
|
private:
|
||||||
|
void tick() {
|
||||||
|
io_context_.poll();
|
||||||
|
}
|
||||||
|
|
||||||
|
std::string read_impl(std::size_t max_size) {
|
||||||
|
tick();
|
||||||
|
if (!port_ || !port_->is_open()) {
|
||||||
|
return "";
|
||||||
|
}
|
||||||
|
auto available = serial::get_available_bytes(port_->native_handle());
|
||||||
|
if (available <= 0) {
|
||||||
|
return "";
|
||||||
|
}
|
||||||
|
auto data = serial::read_all(port_->native_handle());
|
||||||
|
if (data.size() > max_size) {
|
||||||
|
data.resize(max_size);
|
||||||
|
}
|
||||||
|
return data;
|
||||||
|
}
|
||||||
|
|
||||||
|
std::size_t write_impl(const std::string &data) {
|
||||||
|
tick();
|
||||||
|
if (!port_ || !port_->is_open() || data.empty()) {
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
asio::error_code ec;
|
||||||
|
auto size = asio::write(*port_, asio::buffer(data), ec);
|
||||||
|
return ec ? 0 : size;
|
||||||
|
}
|
||||||
|
|
||||||
static std::string normalize_port_name(std::string port_name) {
|
static std::string normalize_port_name(std::string port_name) {
|
||||||
#ifdef _WIN32
|
#ifdef _WIN32
|
||||||
if (port_name.rfind("\\\\.\\", 0) != 0) {
|
if (port_name.rfind("\\\\.\\", 0) != 0) {
|
||||||
|
|||||||
+52
-158
@@ -1,15 +1,14 @@
|
|||||||
#pragma once
|
#pragma once
|
||||||
|
|
||||||
#include "ASIO_Utils.h"
|
#include "ASIO_Utils.h"
|
||||||
#include "Core/Base/Coro_Result.h"
|
|
||||||
#include "TCP_Client.h"
|
#include "TCP_Client.h"
|
||||||
#include "TCP_Server.h"
|
#include "TCP_Server.h"
|
||||||
#include "UDP_Client.h"
|
#include "UDP_Client.h"
|
||||||
#include "UDP_Server.h"
|
#include "UDP_Server.h"
|
||||||
|
|
||||||
#include <concurrencpp/concurrencpp.h>
|
|
||||||
|
|
||||||
#include <array>
|
#include <array>
|
||||||
|
#include <asio/awaitable.hpp>
|
||||||
|
#include <asio/use_awaitable.hpp>
|
||||||
#include <cstddef>
|
#include <cstddef>
|
||||||
#include <cstdint>
|
#include <cstdint>
|
||||||
#include <memory>
|
#include <memory>
|
||||||
@@ -22,27 +21,27 @@ namespace Psc::asio_socket {
|
|||||||
|
|
||||||
class TCP_Client_Coro : public TCP_Client {
|
class TCP_Client_Coro : public TCP_Client {
|
||||||
public:
|
public:
|
||||||
[[nodiscard]] concurrencpp::result<bool> connect_coro() {
|
[[nodiscard]] asio::awaitable<bool> connect_coro() {
|
||||||
auto started = start_connect();
|
auto started = start_connect();
|
||||||
tick();
|
tick();
|
||||||
co_return started;
|
co_return started;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> tick_coro() {
|
[[nodiscard]] asio::awaitable<void> tick_coro() {
|
||||||
tick();
|
tick();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<std::string> read_coro() {
|
[[nodiscard]] asio::awaitable<std::string> read_coro() {
|
||||||
co_return read();
|
co_return read();
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> send_coro(std::string data) {
|
[[nodiscard]] asio::awaitable<void> send_coro(std::string data) {
|
||||||
send(data);
|
send(data);
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> close_coro() {
|
[[nodiscard]] asio::awaitable<void> close_coro() {
|
||||||
close();
|
close();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
@@ -50,42 +49,42 @@ public:
|
|||||||
|
|
||||||
class TCP_Server_Coro : public TCP_Server {
|
class TCP_Server_Coro : public TCP_Server {
|
||||||
public:
|
public:
|
||||||
[[nodiscard]] concurrencpp::result<bool> listen_coro(std::string ip,
|
[[nodiscard]] asio::awaitable<bool> listen_coro(std::string ip,
|
||||||
std::uint32_t port) {
|
std::uint32_t port) {
|
||||||
co_return listen(std::move(ip), port);
|
co_return listen(std::move(ip), port);
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<bool> listen_coro(Sockaddr_In address) {
|
[[nodiscard]] asio::awaitable<bool> listen_coro(Sockaddr_In address) {
|
||||||
co_return listen(address);
|
co_return listen(address);
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> tick_coro() {
|
[[nodiscard]] asio::awaitable<void> tick_coro() {
|
||||||
tick();
|
tick();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> flush_clients_coro() {
|
[[nodiscard]] asio::awaitable<void> flush_clients_coro() {
|
||||||
flush_clients();
|
flush_clients();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void>
|
[[nodiscard]] asio::awaitable<void>
|
||||||
write_to_all_clients_coro(std::string data) {
|
write_to_all_clients_coro(std::string data) {
|
||||||
write_to_all_clients(data);
|
write_to_all_clients(data);
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<std::vector<Read_Info>>
|
[[nodiscard]] asio::awaitable<std::vector<Read_Info>>
|
||||||
read_from_all_clients_coro() {
|
read_from_all_clients_coro() {
|
||||||
co_return read_from_all_clients();
|
co_return read_from_all_clients();
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<std::shared_ptr<TCP_Connect>>
|
[[nodiscard]] asio::awaitable<std::shared_ptr<TCP_Connect>>
|
||||||
accept_coro() const {
|
accept_coro() const {
|
||||||
co_return accept();
|
co_return accept();
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> close_coro() {
|
[[nodiscard]] asio::awaitable<void> close_coro() {
|
||||||
close();
|
close();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
@@ -93,25 +92,25 @@ public:
|
|||||||
|
|
||||||
class UDP_Client_Coro : public UDP_Client {
|
class UDP_Client_Coro : public UDP_Client {
|
||||||
public:
|
public:
|
||||||
[[nodiscard]] concurrencpp::result<bool> connect_coro() {
|
[[nodiscard]] asio::awaitable<bool> connect_coro() {
|
||||||
co_return connect();
|
co_return connect();
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<std::string> read_coro() {
|
[[nodiscard]] asio::awaitable<std::string> read_coro() {
|
||||||
co_return read();
|
co_return read();
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> tick_coro() {
|
[[nodiscard]] asio::awaitable<void> tick_coro() {
|
||||||
tick();
|
tick();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> send_coro(std::string data) {
|
[[nodiscard]] asio::awaitable<void> send_coro(std::string data) {
|
||||||
send(data);
|
send(data);
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> close_coro() {
|
[[nodiscard]] asio::awaitable<void> close_coro() {
|
||||||
close();
|
close();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
@@ -119,36 +118,36 @@ public:
|
|||||||
|
|
||||||
class UDP_Server_Coro : public UDP_Server {
|
class UDP_Server_Coro : public UDP_Server {
|
||||||
public:
|
public:
|
||||||
[[nodiscard]] concurrencpp::result<bool> bind_coro() { co_return bind(); }
|
[[nodiscard]] asio::awaitable<bool> bind_coro() { co_return bind(); }
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> tick_coro() {
|
[[nodiscard]] asio::awaitable<void> tick_coro() {
|
||||||
tick();
|
tick();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<std::vector<Read_Info>> read_coro() {
|
[[nodiscard]] asio::awaitable<std::vector<Read_Info>> read_coro() {
|
||||||
co_return read();
|
co_return read();
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void>
|
[[nodiscard]] asio::awaitable<void>
|
||||||
reply_last_peer_coro(std::string data) {
|
reply_last_peer_coro(std::string data) {
|
||||||
reply_last_peer(data);
|
reply_last_peer(data);
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void>
|
[[nodiscard]] asio::awaitable<void>
|
||||||
send_to_coro(std::string ip, uint16_t port, std::string data) {
|
send_to_coro(std::string ip, uint16_t port, std::string data) {
|
||||||
send_to(ip, port, data);
|
send_to(ip, port, data);
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void>
|
[[nodiscard]] asio::awaitable<void>
|
||||||
write_to_all_clients_coro(std::string msg) {
|
write_to_all_clients_coro(std::string msg) {
|
||||||
write_to_all_clients(msg);
|
write_to_all_clients(msg);
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
[[nodiscard]] concurrencpp::result<void> close_coro() {
|
[[nodiscard]] asio::awaitable<void> close_coro() {
|
||||||
close();
|
close();
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
@@ -173,171 +172,66 @@ inline void throw_if_error(const asio::error_code &ec) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<
|
inline asio::awaitable<
|
||||||
std::shared_ptr<asio::ip::tcp::resolver::results_type>>
|
std::shared_ptr<asio::ip::tcp::resolver::results_type>>
|
||||||
tcp_resolve(asio::ip::tcp::resolver &resolver, std::string host,
|
tcp_resolve(asio::ip::tcp::resolver &resolver, std::string host,
|
||||||
std::uint16_t port) {
|
std::uint16_t port) {
|
||||||
auto endpoints = co_await Psc::coro::callback_result<
|
auto endpoints =
|
||||||
std::shared_ptr<asio::ip::tcp::resolver::results_type>>(
|
co_await resolver.async_resolve(host, std::to_string(port), asio::use_awaitable);
|
||||||
[&resolver, host = std::move(host), port](auto done) mutable {
|
co_return std::make_shared<asio::ip::tcp::resolver::results_type>(std::move(endpoints));
|
||||||
resolver.async_resolve(
|
|
||||||
host, std::to_string(port),
|
|
||||||
[done = std::move(done)](
|
|
||||||
const asio::error_code &ec,
|
|
||||||
asio::ip::tcp::resolver::results_type endpoints) mutable {
|
|
||||||
if (ec) {
|
|
||||||
done.set_exception(
|
|
||||||
std::make_exception_ptr(std::system_error(ec)));
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
done(std::make_shared<asio::ip::tcp::resolver::results_type>(
|
|
||||||
std::move(endpoints)));
|
|
||||||
});
|
|
||||||
});
|
|
||||||
co_return endpoints;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<asio::ip::tcp::endpoint>
|
inline asio::awaitable<asio::ip::tcp::endpoint>
|
||||||
tcp_connect(asio::ip::tcp::socket &socket,
|
tcp_connect(asio::ip::tcp::socket &socket,
|
||||||
const asio::ip::tcp::resolver::results_type &endpoints) {
|
const asio::ip::tcp::resolver::results_type &endpoints) {
|
||||||
struct Result {
|
co_return co_await asio::async_connect(socket, endpoints, asio::use_awaitable);
|
||||||
asio::error_code ec;
|
|
||||||
asio::ip::tcp::endpoint endpoint;
|
|
||||||
};
|
|
||||||
|
|
||||||
auto result = co_await Psc::coro::callback_result<Result>(
|
|
||||||
[&socket, &endpoints](auto done) mutable {
|
|
||||||
asio::async_connect(
|
|
||||||
socket, endpoints,
|
|
||||||
[done = std::move(done)](
|
|
||||||
const asio::error_code &ec,
|
|
||||||
const asio::ip::tcp::endpoint &endpoint) mutable {
|
|
||||||
done(Result{ec, endpoint});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
throw_if_error(result.ec);
|
|
||||||
co_return result.endpoint;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<asio::ip::tcp::endpoint>
|
inline asio::awaitable<asio::ip::tcp::endpoint>
|
||||||
tcp_connect(asio::ip::tcp::socket &socket, asio::ip::tcp::resolver &resolver,
|
tcp_connect(asio::ip::tcp::socket &socket, asio::ip::tcp::resolver &resolver,
|
||||||
std::string host, std::uint16_t port) {
|
std::string host, std::uint16_t port) {
|
||||||
auto endpoints = co_await tcp_resolve(resolver, std::move(host), port);
|
auto endpoints = co_await tcp_resolve(resolver, std::move(host), port);
|
||||||
co_return co_await tcp_connect(socket, *endpoints);
|
co_return co_await tcp_connect(socket, *endpoints);
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<std::shared_ptr<asio::ip::tcp::socket>>
|
inline asio::awaitable<std::shared_ptr<asio::ip::tcp::socket>>
|
||||||
tcp_accept(asio::ip::tcp::acceptor &acceptor) {
|
tcp_accept(asio::ip::tcp::acceptor &acceptor) {
|
||||||
auto socket =
|
auto socket =
|
||||||
std::make_shared<asio::ip::tcp::socket>(acceptor.get_executor());
|
std::make_shared<asio::ip::tcp::socket>(acceptor.get_executor());
|
||||||
|
co_await acceptor.async_accept(*socket, asio::use_awaitable);
|
||||||
auto ec = co_await Psc::coro::callback_result<asio::error_code>(
|
|
||||||
[&acceptor, socket](auto done) mutable {
|
|
||||||
acceptor.async_accept(
|
|
||||||
*socket, [done = std::move(done)](
|
|
||||||
const asio::error_code &ec) mutable { done(ec); });
|
|
||||||
});
|
|
||||||
|
|
||||||
throw_if_error(ec);
|
|
||||||
co_return socket;
|
co_return socket;
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<TCP_Read_Result>
|
inline asio::awaitable<TCP_Read_Result>
|
||||||
tcp_read_some(asio::ip::tcp::socket &socket, std::size_t max_size = 16 * 1024) {
|
tcp_read_some(asio::ip::tcp::socket &socket, std::size_t max_size = 16 * 1024) {
|
||||||
struct Result {
|
std::vector<char> buffer(max_size);
|
||||||
asio::error_code ec;
|
auto size = co_await socket.async_read_some(asio::buffer(buffer), asio::use_awaitable);
|
||||||
std::size_t size{};
|
buffer.resize(size);
|
||||||
};
|
co_return TCP_Read_Result{std::move(buffer)};
|
||||||
|
|
||||||
auto buffer = std::make_shared<std::vector<char>>(max_size);
|
|
||||||
auto result = co_await Psc::coro::callback_result<Result>(
|
|
||||||
[&socket, buffer](auto done) mutable {
|
|
||||||
socket.async_read_some(
|
|
||||||
asio::buffer(*buffer),
|
|
||||||
[done = std::move(done)](const asio::error_code &ec,
|
|
||||||
std::size_t size) mutable {
|
|
||||||
done(Result{ec, size});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
throw_if_error(result.ec);
|
|
||||||
buffer->resize(result.size);
|
|
||||||
co_return TCP_Read_Result{std::move(*buffer)};
|
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<std::size_t>
|
inline asio::awaitable<std::size_t>
|
||||||
tcp_write(asio::ip::tcp::socket &socket, std::string data) {
|
tcp_write(asio::ip::tcp::socket &socket, std::string data) {
|
||||||
struct Result {
|
co_return co_await asio::async_write(socket, asio::buffer(data), asio::use_awaitable);
|
||||||
asio::error_code ec;
|
|
||||||
std::size_t size{};
|
|
||||||
};
|
|
||||||
|
|
||||||
auto buffer = std::make_shared<std::string>(std::move(data));
|
|
||||||
auto result = co_await Psc::coro::callback_result<Result>(
|
|
||||||
[&socket, buffer](auto done) mutable {
|
|
||||||
asio::async_write(socket, asio::buffer(*buffer),
|
|
||||||
[done = std::move(done)](const asio::error_code &ec,
|
|
||||||
std::size_t size) mutable {
|
|
||||||
done(Result{ec, size});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
throw_if_error(result.ec);
|
|
||||||
co_return result.size;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<UDP_Read_Result>
|
inline asio::awaitable<UDP_Read_Result>
|
||||||
udp_receive_from(asio::ip::udp::socket &socket,
|
udp_receive_from(asio::ip::udp::socket &socket,
|
||||||
std::size_t max_size = 16 * 1024) {
|
std::size_t max_size = 16 * 1024) {
|
||||||
struct Result {
|
std::vector<char> buffer(max_size);
|
||||||
asio::error_code ec;
|
asio::ip::udp::endpoint remote;
|
||||||
std::size_t size{};
|
auto size = co_await socket.async_receive_from(asio::buffer(buffer), remote, asio::use_awaitable);
|
||||||
asio::ip::udp::endpoint remote;
|
buffer.resize(size);
|
||||||
};
|
co_return UDP_Read_Result{endpoint_to_sockaddr(remote), std::move(buffer)};
|
||||||
|
|
||||||
auto buffer = std::make_shared<std::vector<char>>(max_size);
|
|
||||||
auto remote = std::make_shared<asio::ip::udp::endpoint>();
|
|
||||||
auto result = co_await Psc::coro::callback_result<Result>(
|
|
||||||
[&socket, buffer, remote](auto done) mutable {
|
|
||||||
socket.async_receive_from(
|
|
||||||
asio::buffer(*buffer), *remote,
|
|
||||||
[done = std::move(done), remote](const asio::error_code &ec,
|
|
||||||
std::size_t size) mutable {
|
|
||||||
done(Result{ec, size, *remote});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
throw_if_error(result.ec);
|
|
||||||
buffer->resize(result.size);
|
|
||||||
co_return UDP_Read_Result{endpoint_to_sockaddr(result.remote),
|
|
||||||
std::move(*buffer)};
|
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<std::size_t>
|
inline asio::awaitable<std::size_t>
|
||||||
udp_send_to(asio::ip::udp::socket &socket, std::string data,
|
udp_send_to(asio::ip::udp::socket &socket, std::string data,
|
||||||
asio::ip::udp::endpoint remote) {
|
asio::ip::udp::endpoint remote) {
|
||||||
struct Result {
|
co_return co_await socket.async_send_to(asio::buffer(data), remote, asio::use_awaitable);
|
||||||
asio::error_code ec;
|
|
||||||
std::size_t size{};
|
|
||||||
};
|
|
||||||
|
|
||||||
auto buffer = std::make_shared<std::string>(std::move(data));
|
|
||||||
auto result = co_await Psc::coro::callback_result<Result>(
|
|
||||||
[&socket, buffer, remote = std::move(remote)](auto done) mutable {
|
|
||||||
socket.async_send_to(
|
|
||||||
asio::buffer(*buffer), remote,
|
|
||||||
[done = std::move(done)](const asio::error_code &ec,
|
|
||||||
std::size_t size) mutable {
|
|
||||||
done(Result{ec, size});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
throw_if_error(result.ec);
|
|
||||||
co_return result.size;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
inline concurrencpp::result<std::size_t>
|
inline asio::awaitable<std::size_t>
|
||||||
udp_send_to(asio::ip::udp::socket &socket, std::string data,
|
udp_send_to(asio::ip::udp::socket &socket, std::string data,
|
||||||
const Sockaddr_In &remote) {
|
const Sockaddr_In &remote) {
|
||||||
asio::error_code ec;
|
asio::error_code ec;
|
||||||
|
|||||||
@@ -195,9 +195,6 @@ if(1)
|
|||||||
|
|
||||||
endif()
|
endif()
|
||||||
|
|
||||||
find_package(concurrencpp CONFIG REQUIRED)
|
|
||||||
|
|
||||||
target_link_libraries(Core_Static PUBLIC concurrencpp::concurrencpp)
|
|
||||||
target_link_libraries(Core_Static PUBLIC boost_pfr)
|
target_link_libraries(Core_Static PUBLIC boost_pfr)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user