/// @file dpf/net/round_sink.hpp /// @brief Transport for round-batched protocol slots. /// @details A session submits one slot per `(round, index)`. The sink writes /// only a contiguous prefix of each round and demuxes the peer's /// matching prefixes into that round's incoming buffer. Memory, /// muxed-trio, and per-round stream sinks all implement this. #ifndef LIBDPF_INCLUDE_DPF_NET_ROUND_SINK_HPP__ #define LIBDPF_INCLUDE_DPF_NET_ROUND_SINK_HPP__ #include #include #include #include #include #include #include #include #include namespace dpf { namespace net { /// @brief Byte transport for one round of a count-sized batch. class RoundSink { public: virtual ~RoundSink() = default; /// @brief How many protocol instances this sink was sized for. virtual std::size_t count() const noexcept = 0; /// @brief How many interactive rounds this sink was sized for. virtual std::size_t rounds() const noexcept = 0; /// @brief Slot width for `round`. Every index of that round uses it. virtual std::size_t slot_bytes(std::uint16_t round) const = 0; /// @brief Store this party's outgoing slot. Index order on the wire is /// enforced at flush time; submits may arrive out of order. virtual void submit(std::uint16_t round, std::size_t index, const std::uint8_t * bytes, std::size_t n) = 0; /// @brief True when the peer's slot for `(round, index)` is available. virtual bool peer_ready(std::uint16_t round, std::size_t index) const = 0; /// @brief Copy the peer's slot into `out`. Requires `peer_ready`. virtual void read_peer(std::uint16_t round, std::size_t index, std::uint8_t * out, std::size_t n) const = 0; /// @brief Send every newly contiguous prefix and pull any peer prefixes. virtual void flush() = 0; /// @brief Flush only `round`'s pending prefix (one barrier on that round). /// @details `sink_exchange` and bit-AND layers call this so a stream sink /// does not touch idle round channels. Default falls back to /// `flush()` for sinks that only implement a global barrier. virtual void flush_round(std::uint16_t /*round*/) { flush(); } /// @brief Service non-blocking I/O. Memory sinks are a no-op. virtual void poll() = 0; /// @brief Block (briefly) until I/O progresses, instead of busy-spinning. /// @details Drive loops call this while waiting on a peer slot. The default /// just `poll()`s and reports "no blocking wait happened" (`false`) /// so existing sinks keep their old behaviour. An event-driven sink /// (e.g. `async_round_sink`) overrides this to sleep in `epoll` via /// `io_context::run_one()`, returning `true` when it ran a handler. /// Returning `false` lets the caller fall back to its spin guard. virtual bool wait_io() { poll(); return false; } /// @brief Wait up to `budget` for one completion. /// @details Memory sinks ignore the budget and behave like `wait_io()`. /// Async sinks sleep in `run_one_for`. virtual bool wait_io_for(std::chrono::milliseconds) { return wait_io(); } /// @brief True when `wait_io_for` sleeps until I/O arrives. Drive loops /// yield instead of sleeping on sinks that cannot block. virtual bool can_block() const noexcept { return false; } /// @brief False while an outbound window is full; pipelined sends pause. virtual bool can_send_ahead() const noexcept { return true; } /// @brief Monotonic count of receive events (peer bytes plus zero-width /// round announcements). Drive loops re-check readiness whenever it /// moves, including when a `poll()` completed the read. virtual std::uint64_t progress() const noexcept { return 0; } /// @brief Sinks that return the same key share one event loop, so blocking /// in one's `wait_io_for` also completes the others' I/O. virtual const void * wait_domain() const noexcept { return this; } }; /// @brief Wait until `sink.peer_ready(round, index)`, for at most `budget`. /// @details Sleeps in the sink's reactor when it can block and yields /// otherwise, so a slow peer costs no CPU and a dead one fails on /// time rather than after a spin count. inline void wait_peer_ready(RoundSink & sink, std::uint16_t round, std::size_t index, std::chrono::milliseconds budget, const char * what) { const auto start = std::chrono::steady_clock::now(); while (!sink.peer_ready(round, index)) { const auto waited = std::chrono::duration_cast( std::chrono::steady_clock::now() - start); if (waited >= budget) throw std::runtime_error(std::string(what) + ": no peer bytes for " + std::to_string(waited.count()) + " ms (budget " + std::to_string(budget.count()) + " ms)"); if (sink.can_block()) sink.wait_io_for(std::min(budget - waited, std::chrono::milliseconds(1000))); else if (!sink.wait_io()) std::this_thread::yield(); } } /// @brief One round's local and peer slot storage with prefix-flush cursors. class round_window { public: round_window() = default; round_window(std::size_t count, std::size_t slot_bytes) : count_(count), slot_bytes_(slot_bytes), written_(count, false), out_(count * slot_bytes), in_(count * slot_bytes) { } std::size_t count() const noexcept { return count_; } std::size_t slot_bytes() const noexcept { return slot_bytes_; } std::size_t next_unwritten() const noexcept { return next_unwritten_; } std::size_t peer_filled() const noexcept { return peer_filled_; } std::size_t flushed() const noexcept { return flushed_; } void submit(std::size_t index, const std::uint8_t * bytes, std::size_t n) { if (index >= count_) throw std::out_of_range("round_window submit index"); if (n != slot_bytes_) throw std::invalid_argument("round_window submit size"); if (written_[index]) throw std::logic_error("round_window double submit"); std::memcpy(out_.data() + index * slot_bytes_, bytes, n); written_[index] = true; while (next_unwritten_ < count_ && written_[next_unwritten_]) ++next_unwritten_; } bool peer_ready(std::size_t index) const noexcept { return index < peer_filled_; } void read_peer(std::size_t index, std::uint8_t * out, std::size_t n) const { if (!peer_ready(index)) throw std::logic_error("round_window peer not ready"); if (n != slot_bytes_) throw std::invalid_argument("round_window read size"); std::memcpy(out, in_.data() + index * slot_bytes_, n); } /// @brief Bytes of the contiguous prefix that have not been flushed yet. const std::uint8_t * pending_out(std::size_t & begin, std::size_t & nslots) const { begin = flushed_; if (next_unwritten_ <= flushed_) { nslots = 0; return nullptr; } nslots = next_unwritten_ - flushed_; return out_.data() + flushed_ * slot_bytes_; } void mark_flushed(std::size_t nslots) { flushed_ += nslots; if (flushed_ > next_unwritten_) throw std::logic_error("round_window flushed past written"); } /// @brief Outbound bytes starting at slot `slot` (retained for resends). const std::uint8_t * out_at(std::size_t slot) const noexcept { return out_.data() + slot * slot_bytes_; } /// @brief Append `nslots` peer slots starting at `peer_filled_`. void accept_peer(const std::uint8_t * bytes, std::size_t nslots) { if (peer_filled_ + nslots > count_) throw std::logic_error("round_window peer overrun"); if (nslots != 0) { std::memcpy(in_.data() + peer_filled_ * slot_bytes_, bytes, nslots * slot_bytes_); } peer_filled_ += nslots; } /// @brief Accept peer slots that begin at an absolute index (mux frames). void accept_peer_at(std::size_t begin, const std::uint8_t * bytes, std::size_t nslots) { if (begin != peer_filled_) throw std::logic_error("round_window peer gap"); accept_peer(bytes, nslots); } private: std::size_t count_ = 0; std::size_t slot_bytes_ = 0; std::size_t next_unwritten_ = 0; std::size_t flushed_ = 0; std::size_t peer_filled_ = 0; std::vector written_; std::vector out_; std::vector in_; }; } // namespace net } // namespace dpf #endif // LIBDPF_INCLUDE_DPF_NET_ROUND_SINK_HPP__