Files
Renderive/web_server/app/Asio_WebSocket_Scene_Loop.h
T
2026-08-18 11:57:43 +08:00

83 lines
3.2 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);
}
auto response = session.handle(event);
if (stopped_.load(std::memory_order_acquire))
co_return;
if (response)
response_handler_(std::move(*response));
}
}
asio::experimental::concurrent_channel<void(asio::error_code, Web_Event)> events_;
std::atomic_bool stopped_{};
Response_Handler response_handler_;
Failure_Handler failure_handler_;
};
}