libdpf/include/dpf/session_host.hpp

95 lines
2.7 KiB
C++
Raw Permalink Normal View History

/// @file dpf/session_host.hpp
/// @brief Durable mesh session that queues micro-plans onto schedule_session.
#ifndef LIBDPF_INCLUDE_DPF_SESSION_HOST_HPP__
#define LIBDPF_INCLUDE_DPF_SESSION_HOST_HPP__
#include <cstddef>
#include <cstdint>
#include <deque>
#include <functional>
#include <map>
#include <stdexcept>
#include <utility>
#include <vector>
#include "dpf/compose.hpp"
#include "dpf/net/edge_mesh.hpp"
#include "dpf/protocol.hpp"
namespace dpf
{
namespace protocol
{
/// @brief Long-lived role on an `edge_mesh` that drives a queue of plans.
/// @details SUBLEQ instructions, hushmap ADD/FIND, and PIRsona GD steps are
/// each one `plan` pushed here. `drive_until_idle` runs
/// `drive_via_schedule` / mesh schedule until every instance is done.
class session_host
{
public:
session_host(std::size_t party, edge_mesh mesh, std::size_t lanes = 1)
: party_(party), mesh_(std::move(mesh)), lanes_(lanes)
{
if (lanes_ == 0)
throw std::invalid_argument("session_host lanes");
}
std::size_t party() const noexcept { return party_; }
edge_mesh & mesh() noexcept { return mesh_; }
std::vector<std::vector<std::uint8_t>> & values() noexcept { return values_; }
void set_kernels(std::map<std::uint32_t, kernel_fn> k)
{
kernels_ = std::move(k);
}
void set_options(drive_options opt) { opt_ = std::move(opt); }
/// @brief Queue a composed plan (micro-plan).
void push(plan p) { queue_.push_back(std::move(p)); }
std::size_t queued() const noexcept { return queue_.size(); }
/// @brief Drive every queued plan to completion on the mesh.
void drive_until_idle()
{
while (!queue_.empty())
{
auto p = std::move(queue_.front());
queue_.pop_front();
drive_one(p);
}
}
/// @brief Drive a single plan without queueing.
void drive_one(const plan & p)
{
if (values_.size() < p.nodes().size())
values_.resize(p.nodes().size());
// Prefer peer-only fast path when the mesh is a single duplex.
if (mesh_.has(edge_peer) && !mesh_.has(edge_rss_next)
&& !mesh_.has(edge_dealer) && mesh_.size() <= 3)
{
drive_via_schedule(p, mesh_.at(edge_peer), values_, kernels_, party_,
opt_);
return;
}
drive_via_schedule(p, mesh_, values_, kernels_, party_, opt_);
}
private:
std::size_t party_ = 0;
edge_mesh mesh_;
std::size_t lanes_ = 1;
std::vector<std::vector<std::uint8_t>> values_;
std::map<std::uint32_t, kernel_fn> kernels_;
drive_options opt_{};
std::deque<plan> queue_;
};
} // namespace protocol
} // namespace dpf
#endif // LIBDPF_INCLUDE_DPF_SESSION_HOST_HPP__