重构协程
This commit is contained in:
@@ -0,0 +1,57 @@
|
||||
#pragma once
|
||||
#include <concurrencpp/concurrencpp.h>
|
||||
#include <exception>
|
||||
#include <memory>
|
||||
#include <type_traits>
|
||||
#include <utility>
|
||||
|
||||
namespace Psc::coro {
|
||||
template <typename T>
|
||||
class Result_Completion {
|
||||
public:
|
||||
explicit Result_Completion(std::shared_ptr<concurrencpp::result_promise<T>> 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<concurrencpp::result_promise<T>> promise_;
|
||||
};
|
||||
template <>
|
||||
class Result_Completion<void> {
|
||||
public:
|
||||
explicit Result_Completion(std::shared_ptr<concurrencpp::result_promise<void>> 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<concurrencpp::result_promise<void>> promise_;
|
||||
};
|
||||
template <typename T, typename Starter>
|
||||
[[nodiscard]] concurrencpp::result<T> callback_result(Starter&& starter)
|
||||
{
|
||||
auto promise = std::make_shared<concurrencpp::result_promise<T>>();
|
||||
auto result = promise->get_result();
|
||||
try {
|
||||
std::forward<Starter>(starter)(Result_Completion<T>{promise});
|
||||
} catch (...) {
|
||||
promise->set_exception(std::current_exception());
|
||||
}
|
||||
return result;
|
||||
}
|
||||
}
|
||||
@@ -5,7 +5,7 @@
|
||||
#include <asio/io_context.hpp>
|
||||
#include <asio/serial_port.hpp>
|
||||
#include <asio/write.hpp>
|
||||
#include <psco/awaitable.hpp>
|
||||
#include <concurrencpp/concurrencpp.h>
|
||||
|
||||
#include <cstddef>
|
||||
#include <memory>
|
||||
@@ -52,13 +52,13 @@ public:
|
||||
io_context_.restart();
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> tick_coro()
|
||||
[[nodiscard]] concurrencpp::result<void> tick_coro()
|
||||
{
|
||||
io_context_.poll();
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<std::string> read_coro(std::size_t max_size = 16 * 1024)
|
||||
[[nodiscard]] concurrencpp::result<std::string> 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<std::size_t> write_coro(std::string data)
|
||||
[[nodiscard]] concurrencpp::result<std::size_t> 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<std::size_t>(size) : static_cast<std::size_t>(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<int64_t>(psco::sync_await(write_coro(data)));
|
||||
return static_cast<int64_t>(write_coro(data).get());
|
||||
}
|
||||
|
||||
int get_available_bytes()
|
||||
|
||||
+45
-47
@@ -5,8 +5,9 @@
|
||||
#include "TCP_Server.h"
|
||||
#include "UDP_Client.h"
|
||||
#include "UDP_Server.h"
|
||||
#include "Core/Base/Coro_Result.h"
|
||||
|
||||
#include <psco/awaitable.hpp>
|
||||
#include <concurrencpp/concurrencpp.h>
|
||||
|
||||
#include <array>
|
||||
#include <cstddef>
|
||||
@@ -21,25 +22,25 @@ namespace Psc::asio_socket {
|
||||
|
||||
class TCP_Client_Coro : public TCP_Client {
|
||||
public:
|
||||
[[nodiscard]] psco::awaitable<bool> connect_coro()
|
||||
[[nodiscard]] concurrencpp::result<bool> connect_coro()
|
||||
{
|
||||
auto started = start_connect();
|
||||
tick();
|
||||
co_return started;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> tick_coro()
|
||||
[[nodiscard]] concurrencpp::result<void> tick_coro()
|
||||
{
|
||||
tick();
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<std::string> read_coro()
|
||||
[[nodiscard]] concurrencpp::result<std::string> read_coro()
|
||||
{
|
||||
co_return read();
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> send_coro(std::string data)
|
||||
[[nodiscard]] concurrencpp::result<void> 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<bool> listen_coro(std::string ip, std::uint32_t port)
|
||||
[[nodiscard]] concurrencpp::result<bool> listen_coro(std::string ip, std::uint32_t port)
|
||||
{
|
||||
co_return listen(std::move(ip), port);
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<bool> listen_coro(Sockaddr_In address)
|
||||
[[nodiscard]] concurrencpp::result<bool> listen_coro(Sockaddr_In address)
|
||||
{
|
||||
co_return listen(address);
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> tick_coro()
|
||||
[[nodiscard]] concurrencpp::result<void> tick_coro()
|
||||
{
|
||||
tick();
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> flush_clients_coro()
|
||||
[[nodiscard]] concurrencpp::result<void> flush_clients_coro()
|
||||
{
|
||||
flush_clients();
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> write_to_all_clients_coro(std::string data)
|
||||
[[nodiscard]] concurrencpp::result<void> write_to_all_clients_coro(std::string data)
|
||||
{
|
||||
write_to_all_clients(data);
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<std::vector<Read_Info>> read_from_all_clients_coro()
|
||||
[[nodiscard]] concurrencpp::result<std::vector<Read_Info>> read_from_all_clients_coro()
|
||||
{
|
||||
co_return read_from_all_clients();
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<std::shared_ptr<TCP_Connect>> accept_coro() const
|
||||
[[nodiscard]] concurrencpp::result<std::shared_ptr<TCP_Connect>> accept_coro() const
|
||||
{
|
||||
co_return accept();
|
||||
}
|
||||
@@ -89,23 +90,23 @@ public:
|
||||
|
||||
class UDP_Client_Coro : public UDP_Client {
|
||||
public:
|
||||
[[nodiscard]] psco::awaitable<bool> connect_coro()
|
||||
[[nodiscard]] concurrencpp::result<bool> connect_coro()
|
||||
{
|
||||
co_return connect();
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<std::string> read_coro()
|
||||
[[nodiscard]] concurrencpp::result<std::string> read_coro()
|
||||
{
|
||||
co_return read();
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> tick_coro()
|
||||
[[nodiscard]] concurrencpp::result<void> tick_coro()
|
||||
{
|
||||
tick();
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> send_coro(std::string data)
|
||||
[[nodiscard]] concurrencpp::result<void> 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<void> tick_coro()
|
||||
[[nodiscard]] concurrencpp::result<void> tick_coro()
|
||||
{
|
||||
tick();
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<std::vector<Read_Info>> read_coro()
|
||||
[[nodiscard]] concurrencpp::result<std::vector<Read_Info>> read_coro()
|
||||
{
|
||||
co_return read();
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> reply_last_peer_coro(std::string data)
|
||||
[[nodiscard]] concurrencpp::result<void> reply_last_peer_coro(std::string data)
|
||||
{
|
||||
reply_last_peer(data);
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> send_to_coro(std::string ip, uint16_t port, std::string data)
|
||||
[[nodiscard]] concurrencpp::result<void> send_to_coro(std::string ip, uint16_t port, std::string data)
|
||||
{
|
||||
send_to(ip, port, data);
|
||||
co_return;
|
||||
}
|
||||
|
||||
[[nodiscard]] psco::awaitable<void> write_to_all_clients_coro(std::string msg)
|
||||
[[nodiscard]] concurrencpp::result<void> 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<asio::ip::tcp::resolver::results_type>
|
||||
inline concurrencpp::result<std::shared_ptr<asio::ip::tcp::resolver::results_type>>
|
||||
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<Result>(
|
||||
auto endpoints = co_await Psc::coro::callback_result<std::shared_ptr<asio::ip::tcp::resolver::results_type>>(
|
||||
[&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<asio::ip::tcp::resolver::results_type>(std::move(endpoints)));
|
||||
});
|
||||
});
|
||||
|
||||
throw_if_error(result.ec);
|
||||
co_return std::move(result.endpoints);
|
||||
co_return endpoints;
|
||||
}
|
||||
|
||||
inline psco::awaitable<asio::ip::tcp::endpoint>
|
||||
inline concurrencpp::result<asio::ip::tcp::endpoint>
|
||||
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<Result>(
|
||||
auto result = co_await Psc::coro::callback_result<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<asio::ip::tcp::endpoint>
|
||||
inline concurrencpp::result<asio::ip::tcp::endpoint>
|
||||
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<std::shared_ptr<asio::ip::tcp::socket>>
|
||||
inline concurrencpp::result<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());
|
||||
|
||||
auto ec = co_await psco::callback_awaitable<asio::error_code>(
|
||||
auto ec = co_await Psc::coro::callback_result<asio::error_code>(
|
||||
[&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<TCP_Read_Result>
|
||||
inline concurrencpp::result<TCP_Read_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<std::vector<char>>(max_size);
|
||||
auto result = co_await psco::callback_awaitable<Result>(
|
||||
auto result = co_await Psc::coro::callback_result<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<std::size_t>
|
||||
inline concurrencpp::result<std::size_t>
|
||||
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::string>(std::move(data));
|
||||
auto result = co_await psco::callback_awaitable<Result>(
|
||||
auto result = co_await Psc::coro::callback_result<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<UDP_Read_Result>
|
||||
inline concurrencpp::result<UDP_Read_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<std::vector<char>>(max_size);
|
||||
auto remote = std::make_shared<asio::ip::udp::endpoint>();
|
||||
auto result = co_await psco::callback_awaitable<Result>(
|
||||
auto result = co_await Psc::coro::callback_result<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<std::size_t>
|
||||
inline concurrencpp::result<std::size_t>
|
||||
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::string>(std::move(data));
|
||||
auto result = co_await psco::callback_awaitable<Result>(
|
||||
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),
|
||||
@@ -345,7 +343,7 @@ udp_send_to(asio::ip::udp::socket& socket,
|
||||
co_return result.size;
|
||||
}
|
||||
|
||||
inline psco::awaitable<std::size_t>
|
||||
inline concurrencpp::result<std::size_t>
|
||||
udp_send_to(asio::ip::udp::socket& socket,
|
||||
std::string data,
|
||||
const Sockaddr_In& remote)
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user