diff --git a/Core/Base/Coro_Result.h b/Core/Base/Coro_Result.h new file mode 100644 index 0000000..8b02b1d --- /dev/null +++ b/Core/Base/Coro_Result.h @@ -0,0 +1,57 @@ +#pragma once +#include +#include +#include +#include +#include + +namespace Psc::coro { +template +class Result_Completion { +public: + explicit Result_Completion(std::shared_ptr> promise) + : promise_(std::move(promise)) + { + } + void operator()(T value) + { + promise_->set_result(std::move(value)); + } + void set_exception(std::exception_ptr exception) + { + promise_->set_exception(exception); + } +private: + std::shared_ptr> promise_; +}; +template <> +class Result_Completion { +public: + explicit Result_Completion(std::shared_ptr> promise) + : promise_(std::move(promise)) + { + } + void operator()() + { + promise_->set_result(); + } + void set_exception(std::exception_ptr exception) + { + promise_->set_exception(exception); + } +private: + std::shared_ptr> promise_; +}; +template +[[nodiscard]] concurrencpp::result callback_result(Starter&& starter) +{ + auto promise = std::make_shared>(); + auto result = promise->get_result(); + try { + std::forward(starter)(Result_Completion{promise}); + } catch (...) { + promise->set_exception(std::current_exception()); + } + return result; +} +} diff --git a/Core/Serial/Serial_Coro.h b/Core/Serial/Serial_Coro.h index 9e34050..8634b8a 100644 --- a/Core/Serial/Serial_Coro.h +++ b/Core/Serial/Serial_Coro.h @@ -5,7 +5,7 @@ #include #include #include -#include +#include #include #include @@ -52,13 +52,13 @@ public: io_context_.restart(); } - [[nodiscard]] psco::awaitable tick_coro() + [[nodiscard]] concurrencpp::result tick_coro() { io_context_.poll(); co_return; } - [[nodiscard]] psco::awaitable read_coro(std::size_t max_size = 16 * 1024) + [[nodiscard]] concurrencpp::result read_coro(std::size_t max_size = 16 * 1024) { co_await tick_coro(); if (!port_ || !port_->is_open()) { @@ -77,7 +77,7 @@ public: co_return data; } - [[nodiscard]] psco::awaitable write_coro(std::string data) + [[nodiscard]] concurrencpp::result write_coro(std::string data) { co_await tick_coro(); if (!port_ || !port_->is_open() || data.empty()) { @@ -96,12 +96,12 @@ public: std::string read(int64_t size = -1) { auto max_size = size > 0 ? static_cast(size) : static_cast(16 * 1024); - return psco::sync_await(read_coro(max_size)); + return (read_coro(max_size)).get(); } int64_t write(const std::string& data) { - return static_cast(psco::sync_await(write_coro(data))); + return static_cast(write_coro(data).get()); } int get_available_bytes() diff --git a/Core/socket/Socket_Coro.h b/Core/socket/Socket_Coro.h index c0c72e4..c9ac8c1 100644 --- a/Core/socket/Socket_Coro.h +++ b/Core/socket/Socket_Coro.h @@ -5,8 +5,9 @@ #include "TCP_Server.h" #include "UDP_Client.h" #include "UDP_Server.h" +#include "Core/Base/Coro_Result.h" -#include +#include #include #include @@ -21,25 +22,25 @@ namespace Psc::asio_socket { class TCP_Client_Coro : public TCP_Client { public: - [[nodiscard]] psco::awaitable connect_coro() + [[nodiscard]] concurrencpp::result connect_coro() { auto started = start_connect(); tick(); co_return started; } - [[nodiscard]] psco::awaitable tick_coro() + [[nodiscard]] concurrencpp::result tick_coro() { tick(); co_return; } - [[nodiscard]] psco::awaitable read_coro() + [[nodiscard]] concurrencpp::result read_coro() { co_return read(); } - [[nodiscard]] psco::awaitable send_coro(std::string data) + [[nodiscard]] concurrencpp::result send_coro(std::string data) { send(data); co_return; @@ -48,40 +49,40 @@ public: class TCP_Server_Coro : public TCP_Server { public: - [[nodiscard]] psco::awaitable listen_coro(std::string ip, std::uint32_t port) + [[nodiscard]] concurrencpp::result listen_coro(std::string ip, std::uint32_t port) { co_return listen(std::move(ip), port); } - [[nodiscard]] psco::awaitable listen_coro(Sockaddr_In address) + [[nodiscard]] concurrencpp::result listen_coro(Sockaddr_In address) { co_return listen(address); } - [[nodiscard]] psco::awaitable tick_coro() + [[nodiscard]] concurrencpp::result tick_coro() { tick(); co_return; } - [[nodiscard]] psco::awaitable flush_clients_coro() + [[nodiscard]] concurrencpp::result flush_clients_coro() { flush_clients(); co_return; } - [[nodiscard]] psco::awaitable write_to_all_clients_coro(std::string data) + [[nodiscard]] concurrencpp::result write_to_all_clients_coro(std::string data) { write_to_all_clients(data); co_return; } - [[nodiscard]] psco::awaitable> read_from_all_clients_coro() + [[nodiscard]] concurrencpp::result> read_from_all_clients_coro() { co_return read_from_all_clients(); } - [[nodiscard]] psco::awaitable> accept_coro() const + [[nodiscard]] concurrencpp::result> accept_coro() const { co_return accept(); } @@ -89,23 +90,23 @@ public: class UDP_Client_Coro : public UDP_Client { public: - [[nodiscard]] psco::awaitable connect_coro() + [[nodiscard]] concurrencpp::result connect_coro() { co_return connect(); } - [[nodiscard]] psco::awaitable read_coro() + [[nodiscard]] concurrencpp::result read_coro() { co_return read(); } - [[nodiscard]] psco::awaitable tick_coro() + [[nodiscard]] concurrencpp::result tick_coro() { tick(); co_return; } - [[nodiscard]] psco::awaitable send_coro(std::string data) + [[nodiscard]] concurrencpp::result send_coro(std::string data) { send(data); co_return; @@ -114,30 +115,30 @@ public: class UDP_Server_Coro : public UDP_Server { public: - [[nodiscard]] psco::awaitable tick_coro() + [[nodiscard]] concurrencpp::result tick_coro() { tick(); co_return; } - [[nodiscard]] psco::awaitable> read_coro() + [[nodiscard]] concurrencpp::result> read_coro() { co_return read(); } - [[nodiscard]] psco::awaitable reply_last_peer_coro(std::string data) + [[nodiscard]] concurrencpp::result reply_last_peer_coro(std::string data) { reply_last_peer(data); co_return; } - [[nodiscard]] psco::awaitable send_to_coro(std::string ip, uint16_t port, std::string data) + [[nodiscard]] concurrencpp::result send_to_coro(std::string ip, uint16_t port, std::string data) { send_to(ip, port, data); co_return; } - [[nodiscard]] psco::awaitable write_to_all_clients_coro(std::string msg) + [[nodiscard]] concurrencpp::result write_to_all_clients_coro(std::string msg) { write_to_all_clients(msg); co_return; @@ -164,32 +165,29 @@ inline void throw_if_error(const asio::error_code& ec) } } -inline psco::awaitable +inline concurrencpp::result> tcp_resolve(asio::ip::tcp::resolver& resolver, std::string host, std::uint16_t port) { - struct Result { - asio::error_code ec; - asio::ip::tcp::resolver::results_type endpoints; - }; - - auto result = co_await psco::callback_awaitable( + auto endpoints = co_await Psc::coro::callback_result>( [&resolver, host = std::move(host), port](auto done) mutable { 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 { - done(Result{ec, std::move(endpoints)}); + if (ec) { + done.set_exception(std::make_exception_ptr(std::system_error(ec))); + return; + } + done(std::make_shared(std::move(endpoints))); }); }); - - throw_if_error(result.ec); - co_return std::move(result.endpoints); + co_return endpoints; } -inline psco::awaitable +inline concurrencpp::result tcp_connect(asio::ip::tcp::socket& socket, const asio::ip::tcp::resolver::results_type& endpoints) { @@ -198,7 +196,7 @@ tcp_connect(asio::ip::tcp::socket& socket, asio::ip::tcp::endpoint endpoint; }; - auto result = co_await psco::callback_awaitable( + auto result = co_await Psc::coro::callback_result( [&socket, &endpoints](auto done) mutable { asio::async_connect( socket, @@ -213,22 +211,22 @@ tcp_connect(asio::ip::tcp::socket& socket, co_return result.endpoint; } -inline psco::awaitable +inline concurrencpp::result tcp_connect(asio::ip::tcp::socket& socket, asio::ip::tcp::resolver& resolver, std::string host, std::uint16_t 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 psco::awaitable> +inline concurrencpp::result> tcp_accept(asio::ip::tcp::acceptor& acceptor) { auto socket = std::make_shared(acceptor.get_executor()); - auto ec = co_await psco::callback_awaitable( + auto ec = co_await Psc::coro::callback_result( [&acceptor, socket](auto done) mutable { acceptor.async_accept( *socket, @@ -241,7 +239,7 @@ tcp_accept(asio::ip::tcp::acceptor& acceptor) co_return socket; } -inline psco::awaitable +inline concurrencpp::result tcp_read_some(asio::ip::tcp::socket& socket, std::size_t max_size = 16 * 1024) { @@ -251,7 +249,7 @@ tcp_read_some(asio::ip::tcp::socket& socket, }; auto buffer = std::make_shared>(max_size); - auto result = co_await psco::callback_awaitable( + auto result = co_await Psc::coro::callback_result( [&socket, buffer](auto done) mutable { socket.async_read_some( asio::buffer(*buffer), @@ -266,7 +264,7 @@ tcp_read_some(asio::ip::tcp::socket& socket, co_return TCP_Read_Result{std::move(*buffer)}; } -inline psco::awaitable +inline concurrencpp::result tcp_write(asio::ip::tcp::socket& socket, std::string data) { @@ -276,7 +274,7 @@ tcp_write(asio::ip::tcp::socket& socket, }; auto buffer = std::make_shared(std::move(data)); - auto result = co_await psco::callback_awaitable( + auto result = co_await Psc::coro::callback_result( [&socket, buffer](auto done) mutable { asio::async_write( socket, @@ -291,7 +289,7 @@ tcp_write(asio::ip::tcp::socket& socket, co_return result.size; } -inline psco::awaitable +inline concurrencpp::result udp_receive_from(asio::ip::udp::socket& socket, std::size_t max_size = 16 * 1024) { @@ -303,7 +301,7 @@ udp_receive_from(asio::ip::udp::socket& socket, auto buffer = std::make_shared>(max_size); auto remote = std::make_shared(); - auto result = co_await psco::callback_awaitable( + auto result = co_await Psc::coro::callback_result( [&socket, buffer, remote](auto done) mutable { socket.async_receive_from( asio::buffer(*buffer), @@ -319,7 +317,7 @@ udp_receive_from(asio::ip::udp::socket& socket, co_return UDP_Read_Result{endpoint_to_sockaddr(result.remote), std::move(*buffer)}; } -inline psco::awaitable +inline concurrencpp::result udp_send_to(asio::ip::udp::socket& socket, std::string data, asio::ip::udp::endpoint remote) @@ -330,7 +328,7 @@ udp_send_to(asio::ip::udp::socket& socket, }; auto buffer = std::make_shared(std::move(data)); - auto result = co_await psco::callback_awaitable( + auto result = co_await Psc::coro::callback_result( [&socket, buffer, remote = std::move(remote)](auto done) mutable { socket.async_send_to( asio::buffer(*buffer), @@ -345,7 +343,7 @@ udp_send_to(asio::ip::udp::socket& socket, co_return result.size; } -inline psco::awaitable +inline concurrencpp::result udp_send_to(asio::ip::udp::socket& socket, std::string data, const Sockaddr_In& remote) diff --git a/main.cmake b/main.cmake index 919ba5c..b8c81d3 100644 --- a/main.cmake +++ b/main.cmake @@ -198,11 +198,6 @@ target_link_libraries(Core_Static PUBLIC concurrencpp::concurrencpp) -set(ucoro_dir "${CMAKE_CURRENT_LIST_DIR}/3rd/psco/src") -file(GLOB_RECURSE ucoro_srcs ${ucoro_dir}/*.h ${ucoro_dir}/*.hpp ${ucoro_dir}/*.cpp) -target_include_directories(Core_Static PUBLIC ${CMAKE_CURRENT_LIST_DIR}/3rd/psco/include) -target_sources(Core_Static PUBLIC ${ucoro_srcs}) - target_link_libraries(Core_Static PUBLIC Core_Interface) target_link_libraries(Core_Static PUBLIC spdlog::spdlog) target_link_libraries(Core_Static PUBLIC asio_object)