573 lines
19 KiB
C++
573 lines
19 KiB
C++
|
|
/// @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 <algorithm>
|
||
|
|
#include <cerrno>
|
||
|
|
#include <chrono>
|
||
|
|
#include <condition_variable>
|
||
|
|
#include <cstddef>
|
||
|
|
#include <cstdint>
|
||
|
|
#include <cstring>
|
||
|
|
#include <fstream>
|
||
|
|
#include <memory>
|
||
|
|
#include <mutex>
|
||
|
|
#include <stdexcept>
|
||
|
|
#include <string>
|
||
|
|
#include <system_error>
|
||
|
|
#include <utility>
|
||
|
|
#include <vector>
|
||
|
|
|
||
|
|
#include <poll.h>
|
||
|
|
#include <sys/socket.h>
|
||
|
|
|
||
|
|
#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<std::vector<std::uint8_t>> a_to_b;
|
||
|
|
std::vector<std::vector<std::uint8_t>> b_to_a;
|
||
|
|
std::vector<std::size_t> a_pos;
|
||
|
|
std::vector<std::size_t> 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<memory_stream_hub> 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<std::mutex> lock(hub_->mu);
|
||
|
|
auto & q = side_a_ ? hub_->a_to_b[i] : hub_->b_to_a[i];
|
||
|
|
const auto * p = static_cast<const std::uint8_t *>(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<std::mutex> 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<memory_stream_hub> hub_;
|
||
|
|
bool side_a_ = true;
|
||
|
|
};
|
||
|
|
|
||
|
|
HEDLEY_WARN_UNUSED_RESULT
|
||
|
|
inline std::pair<memory_stream_array, memory_stream_array>
|
||
|
|
make_memory_stream_pair(std::size_t streams)
|
||
|
|
{
|
||
|
|
auto hub = std::make_shared<memory_stream_hub>(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<const char *>(src),
|
||
|
|
static_cast<std::streamsize>(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<char *>(dst), static_cast<std::streamsize>(n));
|
||
|
|
if (static_cast<std::size_t>(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<std::fstream> 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<std::size_t> 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<std::uint8_t> 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<std::uint8_t>((v >> (8 * i)) & 0xffu);
|
||
|
|
}
|
||
|
|
|
||
|
|
static std::uint32_t get32(const std::uint8_t * p)
|
||
|
|
{
|
||
|
|
return static_cast<std::uint32_t>(p[0]) | (static_cast<std::uint32_t>(p[1]) << 8)
|
||
|
|
| (static_cast<std::uint32_t>(p[2]) << 16)
|
||
|
|
| (static_cast<std::uint32_t>(p[3]) << 24);
|
||
|
|
}
|
||
|
|
|
||
|
|
void send_hello()
|
||
|
|
{
|
||
|
|
const std::size_t r = slot_bytes_.size();
|
||
|
|
std::vector<std::uint8_t> 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<std::uint32_t>(r));
|
||
|
|
put32(h.data() + 12, static_cast<std::uint32_t>(map_.n_lanes));
|
||
|
|
put32(h.data() + 16, static_cast<std::uint32_t>(count_));
|
||
|
|
const std::uint64_t fp = slots_fingerprint(slot_bytes_);
|
||
|
|
for (int i = 0; i < 8; ++i)
|
||
|
|
h[24 + i] = static_cast<std::uint8_t>((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<std::uint64_t>(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<std::uint8_t> 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<std::size_t> slot_bytes_;
|
||
|
|
round_lane_map map_;
|
||
|
|
mutable std::vector<round_window> windows_;
|
||
|
|
std::vector<std::vector<std::uint8_t>> in_;
|
||
|
|
std::vector<std::size_t> in_filled_;
|
||
|
|
std::vector<bool> zero_seen_;
|
||
|
|
std::vector<bool> 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
|