91 lines
4.0 KiB
C++
91 lines
4.0 KiB
C++
#pragma once
|
|
#include "Web_Event.h"
|
|
#include "Web_Plot_Session.h"
|
|
#include <asio/any_io_executor.hpp>
|
|
#include <asio/awaitable.hpp>
|
|
#include <asio/co_spawn.hpp>
|
|
#include <asio/error_code.hpp>
|
|
#include <asio/experimental/concurrent_channel.hpp>
|
|
#include <asio/redirect_error.hpp>
|
|
#include <asio/system_error.hpp>
|
|
#include <asio/use_awaitable.hpp>
|
|
#include <atomic>
|
|
#include <cstddef>
|
|
#include <exception>
|
|
#include <functional>
|
|
#include <memory>
|
|
#include <utility>
|
|
namespace renderive::web {
|
|
template <class Scene_Session>
|
|
class Asio_WebSocket_Scene_Loop final : public std::enable_shared_from_this<Asio_WebSocket_Scene_Loop<Scene_Session>> {
|
|
public:
|
|
enum class Submit_Result {
|
|
none,
|
|
closed,
|
|
queue_full
|
|
};
|
|
using Response_Handler = std::function<void(Web_Response)>;
|
|
using Failure_Handler = std::function<void(std::exception_ptr)>;
|
|
[[nodiscard]] static std::shared_ptr<Asio_WebSocket_Scene_Loop> create(asio::any_io_executor executor, Response_Handler response_handler, Failure_Handler failure_handler) {
|
|
auto loop = std::shared_ptr<Asio_WebSocket_Scene_Loop>(new Asio_WebSocket_Scene_Loop(std::move(executor), std::move(response_handler), std::move(failure_handler)));
|
|
loop->start();
|
|
return loop;
|
|
}
|
|
[[nodiscard]] Submit_Result submit(Web_Event event) {
|
|
if (stopped_.load(std::memory_order_acquire))
|
|
return Submit_Result::closed;
|
|
if (events_.try_send(asio::error_code{}, std::move(event)))
|
|
return Submit_Result::none;
|
|
return stopped_.load(std::memory_order_acquire) ? Submit_Result::closed : Submit_Result::queue_full;
|
|
}
|
|
void stop() {
|
|
if (stopped_.exchange(true, std::memory_order_acq_rel))
|
|
return;
|
|
events_.reset();
|
|
events_.close();
|
|
}
|
|
private:
|
|
static constexpr std::size_t queue_capacity = 128;
|
|
Asio_WebSocket_Scene_Loop(asio::any_io_executor executor, Response_Handler response_handler, Failure_Handler failure_handler)
|
|
: events_(std::move(executor), queue_capacity), response_handler_(std::move(response_handler)), failure_handler_(std::move(failure_handler)) {}
|
|
void start() {
|
|
auto self = this->shared_from_this();
|
|
asio::co_spawn(events_.get_executor(), [self]() -> asio::awaitable<void> {
|
|
co_await self->run();
|
|
}, [self](std::exception_ptr exception) {
|
|
if (exception)
|
|
self->failure_handler_(exception);
|
|
});
|
|
}
|
|
asio::awaitable<void> run() {
|
|
Scene_Session session;
|
|
for (;;) {
|
|
asio::error_code error;
|
|
auto event = co_await events_.async_receive(asio::redirect_error(asio::use_awaitable, error));
|
|
if (error) {
|
|
if (stopped_.load(std::memory_order_acquire))
|
|
co_return;
|
|
throw asio::system_error(error);
|
|
}
|
|
using Completion_Channel = asio::experimental::concurrent_channel<void(asio::error_code, Web_Async_Response)>;
|
|
auto completion_channel = std::make_shared<Completion_Channel>(events_.get_executor(), 1);
|
|
session.async_handle(event, [completion_channel](Web_Async_Response response) {
|
|
completion_channel->try_send(asio::error_code{}, std::move(response));
|
|
});
|
|
asio::error_code completion_error;
|
|
auto completed = co_await completion_channel->async_receive(asio::redirect_error(asio::use_awaitable, completion_error));
|
|
if (completion_error) throw asio::system_error(completion_error);
|
|
if (completed.exception) std::rethrow_exception(completed.exception);
|
|
if (stopped_.load(std::memory_order_acquire))
|
|
co_return;
|
|
if (completed.response)
|
|
response_handler_(std::move(*completed.response));
|
|
}
|
|
}
|
|
asio::experimental::concurrent_channel<void(asio::error_code, Web_Event)> events_;
|
|
std::atomic_bool stopped_{};
|
|
Response_Handler response_handler_;
|
|
Failure_Handler failure_handler_;
|
|
};
|
|
}
|