libdpf/include/dpf/net/sync_stream_array.hpp

112 lines
3.4 KiB
C++
Raw Permalink Normal View History

/// @file dpf/net/sync_stream_array.hpp
/// @brief Blocking `stream_array` face over an `async_stream_array`.
/// @details Setup (TLS, windows, coalescing, reconnect) is the async
/// framework's job. This type only turns each `write` / `read` /
/// `flush` into one async call and pumps that stream's `io_context`
/// until it finishes, so older sync call sites share the same
/// encryption and wire policy as `party_session`.
#ifndef LIBDPF_INCLUDE_DPF_NET_SYNC_STREAM_ARRAY_HPP__
#define LIBDPF_INCLUDE_DPF_NET_SYNC_STREAM_ARRAY_HPP__
#include <atomic>
#include <cstddef>
#include <cstdint>
#include <cstring>
#include <memory>
#include <stdexcept>
#include <string>
#include <system_error>
#include <vector>
#include "dpf/net/async_stream_array.hpp"
#include "dpf/net/buffer_pool.hpp"
#include "dpf/net/stream_array.hpp"
namespace dpf
{
namespace net
{
/// @brief Synchronous view of an `async_stream_array` on one `io_context`.
class sync_stream_array final : public stream_array
{
public:
/// @brief Non-owning: `streams` and its context must outlive this object.
explicit sync_stream_array(async_stream_array & streams)
: streams_(&streams), io_(&streams.context())
{
if (streams_->size() == 0)
throw std::invalid_argument("sync_stream_array empty");
}
std::size_t size() const noexcept override { return streams_->size(); }
void write(std::size_t i, const void * src, std::size_t n) override
{
check(i);
if (n == 0)
return;
auto buf = copy_buffer(src, n);
wait_op([&](auto done) {
streams_->async_write_owned(i, std::move(buf), std::move(done));
});
}
void read(std::size_t i, void * dst, std::size_t n) override
{
check(i);
if (n == 0)
return;
wait_op([&](auto done) {
streams_->async_read(i, dst, n, std::move(done));
});
}
/// @brief No-op: each `write` already waits for the bytes to leave.
void flush(std::size_t) override {}
async_stream_array & async() noexcept { return *streams_; }
const async_stream_array & async() const noexcept { return *streams_; }
private:
void check(std::size_t i) const
{
if (i >= streams_->size())
throw std::out_of_range("sync_stream_array index "
+ std::to_string(i));
}
template <typename Start>
void wait_op(Start start)
{
std::atomic<bool> done{false};
std::error_code ec;
start([&](const std::error_code & e) {
ec = e;
done.store(true, std::memory_order_release);
});
while (!done.load(std::memory_order_acquire))
{
if (io_->stopped())
io_->restart();
io_->run_one();
}
if (ec)
throw std::system_error(ec, "sync_stream_array");
}
async_stream_array * streams_ = nullptr;
asio::io_context * io_ = nullptr;
};
/// @brief Same face as the old fd-based mux: N lanes on one connection.
/// @details Built from an already-established framework link (typically a
/// `party_session` peer). Prefer this over constructing a mux from a
/// raw `channel` socket.
using mux_stream_array = sync_stream_array;
} // namespace net
} // namespace dpf
#endif // LIBDPF_INCLUDE_DPF_NET_SYNC_STREAM_ARRAY_HPP__