module; #include #include #include #include #include #include export module ttwhy.terminal:readers; import :events; import :policies; import :scanner; import ttwhy.core; namespace ttwhy::terminal { export template AppRouter> auto read_events(InputStream & stream, AppRouter & router) -> asio::awaitable { 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{}; queue.reserve(16); auto sink = [&queue](auto const & event) { queue.push_back(event); }; auto scanner = terminal::scanner{sink}; auto raw_buffer = std::array{}; 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{raw_buffer.data(), bytes_read}; scanner.process(byte_span); for (auto const & event : queue) { co_await router.process(event); } queue.clear(); } } } // namespace ttwhy::terminal