/// @file dpf/schedule_streams.hpp /// @brief Run a `schedule_session` over a `stream_array` peer transport. /// @details Builds a multi-round `stream_array_sink` from the schedule's slot /// widths so existing `plan_to_schedule` / compose plans can drive /// without a memory RoundSink. Prefer `drive_plan_on_streams` / /// `drive_both_on_streams` in `compose.hpp` for full plans (they call /// `drive_via_schedule`). Dealer pads still come from a /// `dealer_cursor` or a separate dealer `stream_array`. #ifndef LIBDPF_INCLUDE_DPF_SCHEDULE_STREAMS_HPP__ #define LIBDPF_INCLUDE_DPF_SCHEDULE_STREAMS_HPP__ #include #include #include #include #include #include "hedley/hedley.h" #include "dpf/net/stream_array.hpp" #include "dpf/protocol.hpp" namespace dpf { namespace protocol { /// @brief Slot widths for one stream index per schedule round. HEDLEY_WARN_UNUSED_RESULT inline std::vector slot_bytes_of( const std::vector & rounds) { std::vector out; out.reserve(rounds.size()); for (const auto & r : rounds) out.push_back(r.slot_bytes); return out; } /// @brief Owning peer sink + schedule session for stream-array transport. /// @details Requires `streams.size() >= rounds.size()`. struct owning_schedule_on_streams { std::unique_ptr sink; schedule_session session; owning_schedule_on_streams(std::size_t count, net::stream_array & streams, std::vector rounds) : sink(std::make_unique(streams, slot_bytes_of(rounds), count)), session(count, *sink, std::move(rounds)) { } void submit(std::size_t index) { session.submit(index); } void drive() { session.drive(); } }; HEDLEY_WARN_UNUSED_RESULT inline owning_schedule_on_streams make_owning_schedule_on_streams( std::size_t count, net::stream_array & streams, std::vector rounds) { return owning_schedule_on_streams(count, streams, std::move(rounds)); } } // namespace protocol } // namespace dpf #endif