141 lines
4 KiB
C++
141 lines
4 KiB
C++
|
|
/// @file dpf/net/tcp_mesh.hpp
|
||
|
|
/// @brief N-party TCP clique bootstrap via `party_session`.
|
||
|
|
/// @details Thin wrappers around `party_session::join` so sync and async mesh
|
||
|
|
/// helpers share TLS, wire policy, deadlines, and reconnect with the
|
||
|
|
/// rest of the framework. Prefer constructing a `party_session`
|
||
|
|
/// directly in new code.
|
||
|
|
#ifndef LIBDPF_INCLUDE_DPF_NET_TCP_MESH_HPP__
|
||
|
|
#define LIBDPF_INCLUDE_DPF_NET_TCP_MESH_HPP__
|
||
|
|
|
||
|
|
#include <cstddef>
|
||
|
|
#include <memory>
|
||
|
|
#include <stdexcept>
|
||
|
|
#include <string>
|
||
|
|
#include <utility>
|
||
|
|
#include <vector>
|
||
|
|
|
||
|
|
#include "dpf/net/asio_ns.hpp"
|
||
|
|
|
||
|
|
#include "dpf/net/mesh_rendezvous.hpp"
|
||
|
|
#include "dpf/net/party_session.hpp"
|
||
|
|
#include "dpf/net/policy.hpp"
|
||
|
|
#include "dpf/net/security.hpp"
|
||
|
|
#include "dpf/net/sync_stream_array.hpp"
|
||
|
|
|
||
|
|
namespace dpf
|
||
|
|
{
|
||
|
|
namespace net
|
||
|
|
{
|
||
|
|
|
||
|
|
/// @brief One party's view of a TCP mux clique backed by `party_session`.
|
||
|
|
struct async_tcp_mesh
|
||
|
|
{
|
||
|
|
unsigned me = 0;
|
||
|
|
unsigned parties = 0;
|
||
|
|
std::unique_ptr<party_session> session;
|
||
|
|
|
||
|
|
async_stream_array & peer(unsigned p)
|
||
|
|
{
|
||
|
|
if (!session)
|
||
|
|
throw std::logic_error("async_tcp_mesh: not joined");
|
||
|
|
return session->peer(p);
|
||
|
|
}
|
||
|
|
|
||
|
|
bool has(unsigned p) const
|
||
|
|
{
|
||
|
|
return session && p < parties && p != me;
|
||
|
|
}
|
||
|
|
|
||
|
|
party_session & sess()
|
||
|
|
{
|
||
|
|
if (!session)
|
||
|
|
throw std::logic_error("async_tcp_mesh: not joined");
|
||
|
|
return *session;
|
||
|
|
}
|
||
|
|
};
|
||
|
|
|
||
|
|
/// @brief Sync mux clique view (same topology as `async_tcp_mesh`).
|
||
|
|
/// @details Owns the `io_context` and a sync facade per peer edge.
|
||
|
|
struct tcp_mesh
|
||
|
|
{
|
||
|
|
unsigned me = 0;
|
||
|
|
unsigned parties = 0;
|
||
|
|
std::shared_ptr<asio::io_context> io;
|
||
|
|
std::unique_ptr<party_session> session;
|
||
|
|
std::vector<std::unique_ptr<mux_stream_array>> links;
|
||
|
|
|
||
|
|
mux_stream_array & peer(unsigned p)
|
||
|
|
{
|
||
|
|
if (p == me || p >= parties || !links[p])
|
||
|
|
throw std::invalid_argument("tcp_mesh: bad peer");
|
||
|
|
return *links[p];
|
||
|
|
}
|
||
|
|
|
||
|
|
party_session & sess()
|
||
|
|
{
|
||
|
|
if (!session)
|
||
|
|
throw std::logic_error("tcp_mesh: not joined");
|
||
|
|
return *session;
|
||
|
|
}
|
||
|
|
};
|
||
|
|
|
||
|
|
inline session_options mesh_session_options(std::size_t nstreams,
|
||
|
|
const wire_policy & pol = {}, const deadlines & lim = {},
|
||
|
|
const peer_security & sec = {})
|
||
|
|
{
|
||
|
|
session_options opt;
|
||
|
|
opt.n_lanes = nstreams;
|
||
|
|
opt.kind = transport::mux;
|
||
|
|
opt.policy = pol;
|
||
|
|
opt.limits = lim;
|
||
|
|
opt.security = sec;
|
||
|
|
return opt;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Party `me` joins an N-party async mux clique on `host`.
|
||
|
|
inline async_tcp_mesh join_async_tcp_mesh(asio::io_context & io, unsigned me,
|
||
|
|
unsigned n, const std::string & host, mesh_ports & ports,
|
||
|
|
std::size_t nstreams, const wire_policy & pol = {},
|
||
|
|
const deadlines & lim = {}, const peer_security & sec = {})
|
||
|
|
{
|
||
|
|
if (nstreams == 0)
|
||
|
|
throw std::invalid_argument("join_async_tcp_mesh: nstreams");
|
||
|
|
async_tcp_mesh mesh;
|
||
|
|
mesh.me = me;
|
||
|
|
mesh.parties = n;
|
||
|
|
mesh.session = std::make_unique<party_session>(io, me, n,
|
||
|
|
mesh_session_options(nstreams, pol, lim, sec));
|
||
|
|
mesh.session->join(host, ports, nstreams);
|
||
|
|
return mesh;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Sync variant: each edge is a `mux_stream_array` over the session.
|
||
|
|
inline tcp_mesh join_tcp_mesh(unsigned me, unsigned n, const std::string & host,
|
||
|
|
mesh_ports & ports, std::size_t nstreams, const wire_policy & pol = {},
|
||
|
|
const deadlines & lim = {}, const peer_security & sec = {})
|
||
|
|
{
|
||
|
|
if (nstreams == 0)
|
||
|
|
throw std::invalid_argument("join_tcp_mesh: nstreams");
|
||
|
|
tcp_mesh mesh;
|
||
|
|
mesh.me = me;
|
||
|
|
mesh.parties = n;
|
||
|
|
mesh.io = std::make_shared<asio::io_context>();
|
||
|
|
mesh.session = std::make_unique<party_session>(*mesh.io, me, n,
|
||
|
|
mesh_session_options(nstreams, pol, lim, sec));
|
||
|
|
mesh.session->join(host, ports, nstreams);
|
||
|
|
mesh.links.resize(n);
|
||
|
|
for (unsigned peer = 0; peer < n; ++peer)
|
||
|
|
{
|
||
|
|
if (peer == me)
|
||
|
|
continue;
|
||
|
|
mesh.links[peer] =
|
||
|
|
std::make_unique<mux_stream_array>(mesh.session->peer(peer));
|
||
|
|
}
|
||
|
|
return mesh;
|
||
|
|
}
|
||
|
|
|
||
|
|
} // namespace net
|
||
|
|
} // namespace dpf
|
||
|
|
|
||
|
|
#endif // LIBDPF_INCLUDE_DPF_NET_TCP_MESH_HPP__
|