aboutsummaryrefslogtreecommitdiff
path: root/ttwhy/io/readers.cppm
diff options
context:
space:
mode:
Diffstat (limited to 'ttwhy/io/readers.cppm')
-rw-r--r--ttwhy/io/readers.cppm93
1 files changed, 93 insertions, 0 deletions
diff --git a/ttwhy/io/readers.cppm b/ttwhy/io/readers.cppm
new file mode 100644
index 0000000..d15cc7a
--- /dev/null
+++ b/ttwhy/io/readers.cppm
@@ -0,0 +1,93 @@
+module;
+
+#include <asio.hpp>
+#include <asio/experimental/awaitable_operators.hpp>
+
+#include <array>
+#include <chrono>
+#include <span>
+#include <vector>
+
+export module ttwhy.io:readers;
+
+import ttwhy.routers;
+import ttwhy.scanners;
+
+namespace ttwhy::io
+{
+
+ export template<typename InputStream, router AppRouter>
+ auto read_events(InputStream & stream, AppRouter & router) -> asio::awaitable<void>
+ {
+ using namespace asio::experimental::awaitable_operators;
+ using namespace std::chrono_literals;
+
+ auto executor = co_await asio::this_coro::executor;
+ auto timer = asio::steady_timer{executor};
+
+ auto queue = std::vector<scanners::input_event>{};
+ queue.reserve(16);
+
+ auto sink = [&queue](auto const & event) {
+ queue.push_back(event);
+ };
+
+ using terminal_policy = ttwhy::scanners::associated_terminal_policy_t<AppRouter>;
+ auto scanner = scanners::terminal_scanner<decltype(sink), terminal_policy>{sink};
+
+ auto raw_buffer = std::array<char, 64>{};
+
+ while (true)
+ {
+ auto error = asio::error_code{};
+ auto bytes_read = 0uz;
+
+ if (scanner.is_pending())
+ {
+ timer.expires_after(50ms);
+
+ auto result = co_await (stream.async_read_some(asio::buffer(raw_buffer), asio::as_tuple(asio::use_awaitable)) ||
+ timer.async_wait(asio::as_tuple(asio::use_awaitable)));
+
+ if (result.index() == 0)
+ {
+ std::tie(error, bytes_read) = std::get<0>(result);
+ }
+ else
+ {
+ scanner.timeout();
+ for (auto const & event : queue)
+ {
+ co_await router.process(event);
+ }
+ queue.clear();
+ continue;
+ }
+ }
+ else
+ {
+ std::tie(error, bytes_read) =
+ co_await stream.async_read_some(asio::buffer(raw_buffer), asio::as_tuple(asio::use_awaitable));
+ }
+
+ if (error)
+ {
+ if (error == asio::error::interrupted)
+ {
+ continue;
+ }
+ co_return;
+ }
+
+ auto const byte_span = std::span<char const>{raw_buffer.data(), bytes_read};
+ scanner.process(byte_span);
+
+ for (auto const & event : queue)
+ {
+ co_await router.process(event);
+ }
+ queue.clear();
+ }
+ }
+
+} // namespace ttwhy::io