/// @file dpf/net/stream_array.hpp /// @brief Indexed pack of duplex byte streams for dealer and online rounds. /// @details `read` / `write` / `flush` are the only operations. Backends are /// memory (tests), file (stored prep), and one-channel mux (TCP). /// SCTP is a later backend: index `i` maps to association stream `i`. /// Protocol code must not call TCP or SCTP sockets directly. #ifndef LIBDPF_INCLUDE_DPF_NET_STREAM_ARRAY_HPP__ #define LIBDPF_INCLUDE_DPF_NET_STREAM_ARRAY_HPP__ #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include "hedley/hedley.h" #include "dpf/net/channel.hpp" #include "dpf/net/policy.hpp" #include "dpf/net/round_lane.hpp" #include "dpf/net/round_sink.hpp" namespace dpf { namespace net { /// @brief Virtual base for N duplex streams. class stream_array { public: virtual ~stream_array() = default; virtual std::size_t size() const noexcept = 0; virtual void write(std::size_t i, const void * src, std::size_t n) = 0; virtual void read(std::size_t i, void * dst, std::size_t n) = 0; virtual void flush(std::size_t i) = 0; void write_all(std::size_t i, const void * src, std::size_t n) { write(i, src, n); flush(i); } }; // --------------------------------------------------------------------------- // Memory // --------------------------------------------------------------------------- struct memory_stream_hub { std::size_t n = 0; std::vector> a_to_b; std::vector> b_to_a; std::vector a_pos; std::vector b_pos; mutable std::mutex mu; std::condition_variable cv; explicit memory_stream_hub(std::size_t streams) : n(streams), a_to_b(streams), b_to_a(streams), a_pos(streams, 0), b_pos(streams, 0) { } }; /// @brief One end of a paired in-process stream array. class memory_stream_array : public stream_array { public: memory_stream_array(std::shared_ptr hub, bool side_a) : hub_(std::move(hub)), side_a_(side_a) { if (!hub_) throw std::invalid_argument("memory_stream_array needs a hub"); } std::size_t size() const noexcept override { return hub_->n; } void write(std::size_t i, const void * src, std::size_t n) override { check(i); { std::lock_guard lock(hub_->mu); auto & q = side_a_ ? hub_->a_to_b[i] : hub_->b_to_a[i]; const auto * p = static_cast(src); q.insert(q.end(), p, p + n); } hub_->cv.notify_all(); } void read(std::size_t i, void * dst, std::size_t n) override { check(i); std::unique_lock lock(hub_->mu); auto & q = side_a_ ? hub_->b_to_a[i] : hub_->a_to_b[i]; auto & pos = side_a_ ? hub_->b_pos[i] : hub_->a_pos[i]; hub_->cv.wait(lock, [&] { return pos + n <= q.size(); }); std::memcpy(dst, q.data() + pos, n); pos += n; } void flush(std::size_t) override {} private: void check(std::size_t i) const { if (i >= hub_->n) throw std::out_of_range("memory_stream_array index"); } std::shared_ptr hub_; bool side_a_ = true; }; HEDLEY_WARN_UNUSED_RESULT inline std::pair make_memory_stream_pair(std::size_t streams) { auto hub = std::make_shared(streams); return {memory_stream_array(hub, true), memory_stream_array(hub, false)}; } // --------------------------------------------------------------------------- // File (one file per stream index) // --------------------------------------------------------------------------- /// @brief Read-only or write-only file per stream. Dealer writes; online reads. class file_stream_array : public stream_array { public: /// @brief Open `basename + "-" + i` for each index. /// @param write_mode true for dealer output, false for online input. file_stream_array(const std::string & basename, std::size_t streams, bool write_mode) : write_mode_(write_mode), files_(streams) { for (std::size_t i = 0; i < streams; ++i) { const auto path = basename + "-" + std::to_string(i); files_[i].open(path, write_mode ? (std::ios::binary | std::ios::trunc | std::ios::out) : (std::ios::binary | std::ios::in)); if (!files_[i]) throw std::runtime_error("file_stream_array: " + path); } } std::size_t size() const noexcept override { return files_.size(); } void write(std::size_t i, const void * src, std::size_t n) override { if (!write_mode_) throw std::logic_error("file_stream_array: write on read-only"); check(i); files_[i].write(static_cast(src), static_cast(n)); if (!files_[i]) throw std::runtime_error("file_stream_array: write failed"); } void read(std::size_t i, void * dst, std::size_t n) override { if (write_mode_) throw std::logic_error("file_stream_array: read on write-only"); check(i); files_[i].read(static_cast(dst), static_cast(n)); if (static_cast(files_[i].gcount()) != n) throw std::runtime_error("file_stream_array: short read"); } void flush(std::size_t i) override { check(i); if (write_mode_) files_[i].flush(); } private: void check(std::size_t i) const { if (i >= files_.size()) throw std::out_of_range("file_stream_array index"); } bool write_mode_ = false; std::vector files_; }; // --------------------------------------------------------------------------- // Mux wire header (shared with async_mux_stream_array) // --------------------------------------------------------------------------- #pragma pack(push, 1) struct stream_frame_hdr { std::uint16_t index = 0; std::uint32_t length = 0; }; #pragma pack(pop) // Synchronous mux is `sync_stream_array` / `mux_stream_array` in // dpf/net/sync_stream_array.hpp: a blocking face over the async mux backends // (including TLS), so sync call sites share party_session encryption and // wire policy. // --------------------------------------------------------------------------- // RoundSink adapter (round → stream, with lane reuse when rounds > streams) // --------------------------------------------------------------------------- /// @brief Multi-round RoundSink over a synchronous `stream_array`. /// @details Same wire contract as `async_round_sink`: a plan-shape hello on /// lane 0 (sent at construction, read before the first peer bytes), /// unframed lanes when every round has its own lane, and framed /// `{round, nbytes}` lanes otherwise (or when `framing` says so). A /// framed round accepts partial instance prefixes. Call /// `flush_round(r)` before starting another round's pending writes. class stream_array_sink : public RoundSink { public: stream_array_sink(stream_array & streams, std::vector slot_bytes_per_round, std::size_t count = 1, framing_mode framing = framing_mode::automatic, bool hello = true) : streams_(&streams), count_(count), slot_bytes_(std::move(slot_bytes_per_round)), map_(streams.size(), slot_bytes_.size(), framing), windows_(slot_bytes_.size()), in_(slot_bytes_.size()), in_filled_(slot_bytes_.size(), 0), zero_seen_(slot_bytes_.size(), false), announced_(slot_bytes_.size(), false), hello_(hello) { if (streams_->size() == 0) throw std::invalid_argument("stream_array_sink: empty streams"); if (count_ == 0) throw std::invalid_argument("stream_array_sink: count 0"); for (std::size_t r = 0; r < slot_bytes_.size(); ++r) { windows_[r] = round_window(count_, slot_bytes_[r]); in_[r].assign(count_ * slot_bytes_[r], 0); } if (hello_) send_hello(); } /// @brief Legacy single-round sink on stream `index` (no hello). stream_array_sink(stream_array & streams, std::size_t index, std::size_t slot) : streams_(&streams), count_(1), legacy_index_(index), use_legacy_index_(true), map_(1, 1) { slot_bytes_.push_back(slot); windows_.emplace_back(count_, slot); in_.emplace_back(count_ * slot, 0); in_filled_.push_back(0); zero_seen_.push_back(false); announced_.push_back(false); } std::size_t count() const noexcept override { return count_; } std::size_t rounds() const noexcept override { return slot_bytes_.size(); } bool framed() const noexcept { return map_.framed && !use_legacy_index_; } std::size_t slot_bytes(std::uint16_t round) const override { if (round >= slot_bytes_.size()) throw std::out_of_range("stream_array_sink round"); return slot_bytes_[round]; } void submit(std::uint16_t round, std::size_t index, const std::uint8_t * bytes, std::size_t n) override { window(round).submit(index, bytes, n); } bool peer_ready(std::uint16_t round, std::size_t index) const override { if (round >= windows_.size()) return false; if (framed()) { if (slot_bytes_[round] == 0) return zero_seen_[round]; return in_filled_[round] >= (index + 1) * slot_bytes_[round]; } return window(round).peer_ready(index); } void read_peer(std::uint16_t round, std::size_t index, std::uint8_t * out, std::size_t n) const override { if (framed()) { if (!peer_ready(round, index)) throw std::logic_error("stream_array_sink: peer not ready"); if (n != slot_bytes_[round]) throw std::invalid_argument("stream_array_sink read size"); if (n != 0) std::memcpy(out, in_[round].data() + index * n, n); return; } window(round).read_peer(index, out, n); } void flush() override { for (std::uint16_t r = 0; r < slot_bytes_.size(); ++r) flush_round(r); } void flush_round(std::uint16_t round) override { auto & win = window(round); std::size_t begin = 0; std::size_t nslots = 0; const auto * pend = win.pending_out(begin, nslots); const std::size_t sb = slot_bytes_[round]; const std::size_t stream_i = stream_index(round); const std::size_t nbytes = nslots * sb; if (framed()) { const bool announce = sb == 0 && !announced_[round]; if (nslots == 0 && !announce) return; auto frame = pack_round_frame(round, pend, nbytes); streams_->write(stream_i, frame.data(), frame.size()); streams_->flush(stream_i); if (announce) announced_[round] = true; if (nslots != 0) win.mark_flushed(nslots); ensure_hello(); while (sb == 0 ? !zero_seen_[round] : in_filled_[round] < win.flushed() * sb) pull_one_frame(stream_i); return; } if (nslots != 0) { streams_->write(stream_i, pend, nbytes); streams_->flush(stream_i); win.mark_flushed(nslots); } if (nslots == 0) return; ensure_hello(); std::vector peer(nbytes); streams_->read(stream_i, peer.data(), peer.size()); win.accept_peer_at(begin, peer.data(), nslots); } void poll() override {} private: static constexpr std::uint32_t k_magic = 0x48535044u; static constexpr std::size_t k_fixed = 36; std::size_t stream_index(std::uint16_t round) const { if (use_legacy_index_) return legacy_index_; return map_.lane(round); } static void put32(std::uint8_t * p, std::uint32_t v) { for (int i = 0; i < 4; ++i) p[i] = static_cast((v >> (8 * i)) & 0xffu); } static std::uint32_t get32(const std::uint8_t * p) { return static_cast(p[0]) | (static_cast(p[1]) << 8) | (static_cast(p[2]) << 16) | (static_cast(p[3]) << 24); } void send_hello() { const std::size_t r = slot_bytes_.size(); std::vector h(k_fixed + 4 * r, 0); put32(h.data(), k_magic); h[4] = 1; h[6] = map_.framed ? 1 : 0; put32(h.data() + 8, static_cast(r)); put32(h.data() + 12, static_cast(map_.n_lanes)); put32(h.data() + 16, static_cast(count_)); const std::uint64_t fp = slots_fingerprint(slot_bytes_); for (int i = 0; i < 8; ++i) h[24 + i] = static_cast((fp >> (8 * i)) & 0xffu); streams_->write(0, h.data(), h.size()); streams_->flush(0); } void ensure_hello() { if (!hello_ || hello_read_) return; hello_read_ = true; std::uint8_t f[k_fixed]; streams_->read(0, f, k_fixed); if (get32(f) != k_magic) throw std::runtime_error("stream_array_sink: peer did not send a " "sink hello (both ends must use sinks with hello enabled)"); const std::uint32_t pr = get32(f + 8); const std::uint32_t pl = get32(f + 12); const std::uint32_t pc = get32(f + 16); const bool pfr = (f[6] & 1) != 0; std::uint64_t pf = 0; for (int i = 0; i < 8; ++i) pf |= static_cast(f[24 + i]) << (8 * i); std::string why; if (pr != slot_bytes_.size()) why += " rounds " + std::to_string(pr) + " vs " + std::to_string(slot_bytes_.size()) + ";"; else if (pf != slots_fingerprint(slot_bytes_)) why += " slot widths differ;"; if (pl != map_.n_lanes) why += " lanes " + std::to_string(pl) + " vs " + std::to_string(map_.n_lanes) + ";"; if (pc != count_) why += " instances " + std::to_string(pc) + " vs " + std::to_string(count_) + ";"; if (pfr != map_.framed) why += std::string(" framing ") + (pfr ? "on" : "off") + " vs " + (map_.framed ? "on" : "off") + ";"; if (!why.empty()) throw std::runtime_error( "stream_array_sink: peer sink disagrees (peer vs this side):" + why); std::vector rest(4 * pr); if (!rest.empty()) streams_->read(0, rest.data(), rest.size()); } void pull_one_frame(std::size_t stream_i) { std::uint8_t raw[round_lane_hdr::size]; streams_->read(stream_i, raw, round_lane_hdr::size); const auto hdr = round_lane_hdr::unpack(raw); if (hdr.round >= slot_bytes_.size() || map_.lane(hdr.round) != stream_i) throw std::runtime_error("stream_array_sink: frame for round " + std::to_string(hdr.round) + " on lane " + std::to_string(stream_i)); const std::size_t sb = slot_bytes_[hdr.round]; if (hdr.nbytes == 0) { if (sb == 0) zero_seen_[hdr.round] = true; return; } if (sb == 0 || hdr.nbytes % sb != 0 || in_filled_[hdr.round] + hdr.nbytes > count_ * sb) throw std::runtime_error("stream_array_sink: round " + std::to_string(hdr.round) + " frame of " + std::to_string(hdr.nbytes) + " bytes does not fit"); streams_->read(stream_i, in_[hdr.round].data() + in_filled_[hdr.round], hdr.nbytes); in_filled_[hdr.round] += hdr.nbytes; } round_window & window(std::uint16_t round) { if (round >= windows_.size()) throw std::out_of_range("stream_array_sink round"); return windows_[round]; } const round_window & window(std::uint16_t round) const { if (round >= windows_.size()) throw std::out_of_range("stream_array_sink round"); return windows_[round]; } stream_array * streams_ = nullptr; std::size_t count_ = 1; std::size_t legacy_index_ = 0; bool use_legacy_index_ = false; std::vector slot_bytes_; round_lane_map map_; mutable std::vector windows_; std::vector> in_; std::vector in_filled_; std::vector zero_seen_; std::vector announced_; bool hello_ = false; bool hello_read_ = false; }; /// @brief Map index `i` to SCTP stream `i` on one association. /// @details Synchronous stub: this blocking `stream_array` face has no linked OS /// backend, so its I/O throws. Protocol code should depend on /// `stream_array` only; swap in a linked backend without API changes. /// The real, event-driven SCTP backend already exists on Linux as /// `dpf::net::async_sctp_stream_array` (see /// `net/async_sctp_stream_array.hpp`, guarded by `DPF_HAS_LIBSCTP`). /// It implements `dpf::net::async_stream_array`, still maps index `i` /// to SCTP stream `i`, and drops in behind /// `dpf::async::overlapped_byte_protocol` with no protocol changes. /// Link `-lsctp`. Windows has no supported SCTP path (SctpDrv is /// unstable), so the async class throws there. struct sctp_association { int fd = -1; }; class sctp_stream_array : public stream_array { public: /// @brief Reserved for a linked backend (`libsctp` association fd). explicit sctp_stream_array(sctp_association /*assoc*/) { throw std::logic_error("sctp_stream_array: not linked"); } /// @brief Size-only stub: every I/O op throws until a backend is linked. explicit sctp_stream_array(std::size_t streams) : n_(streams) { if (streams == 0) throw std::invalid_argument("sctp_stream_array empty"); } std::size_t size() const noexcept override { return n_; } void write(std::size_t, const void *, std::size_t) override { throw std::logic_error( "sctp_stream_array: write — SCTP backend not linked (index i -> SCTP stream i)"); } void read(std::size_t, void *, std::size_t) override { throw std::logic_error( "sctp_stream_array: read — SCTP backend not linked (index i -> SCTP stream i)"); } void flush(std::size_t) override { throw std::logic_error( "sctp_stream_array: flush — SCTP backend not linked (index i -> SCTP stream i)"); } private: std::size_t n_ = 0; }; } // namespace net } // namespace dpf #endif