diff options
| author | Felix Morgner <felix.morgner@gmail.com> | 2026-09-22 22:47:33 +0200 |
|---|---|---|
| committer | Felix Morgner <felix.morgner@gmail.com> | 2026-09-22 22:47:33 +0200 |
| commit | d05b9343f69a0841043963088e0003368e205bd9 (patch) | |
| tree | 8cde07ff86d27cbc808a0a149b89d0024b4a586c /ttwhy/terminal/readers.cppm | |
| parent | 22f6534081acfabdab4956d337604a21e7da5b64 (diff) | |
| download | ttwhy-d05b9343f69a0841043963088e0003368e205bd9.tar.xz ttwhy-d05b9343f69a0841043963088e0003368e205bd9.zip | |
Diffstat (limited to 'ttwhy/terminal/readers.cppm')
| -rw-r--r-- | ttwhy/terminal/readers.cppm | 94 |
1 files changed, 94 insertions, 0 deletions
diff --git a/ttwhy/terminal/readers.cppm b/ttwhy/terminal/readers.cppm new file mode 100644 index 0000000..ffa5687 --- /dev/null +++ b/ttwhy/terminal/readers.cppm @@ -0,0 +1,94 @@ +module; + +#include <asio.hpp> +#include <asio/experimental/awaitable_operators.hpp> + +#include <array> +#include <chrono> +#include <span> +#include <vector> + +export module ttwhy.terminal:readers; + +import :events; +import :policies; +import :scanner; +import ttwhy.core; + +namespace ttwhy::terminal +{ + + export template<typename TerminalPolicy = xterm_policy, typename InputStream, ttwhy::router<input_event> 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<input_event>{}; + queue.reserve(16); + + auto sink = [&queue](auto const & event) { + queue.push_back(event); + }; + + auto scanner = terminal::scanner<decltype(sink), TerminalPolicy>{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::terminal |
