877 lines
31 KiB
C++
877 lines
31 KiB
C++
|
|
/// @file dpf/experiment.hpp
|
||
|
|
/// @brief Replayable master seed + paper-ready protocol cost CSVs.
|
||
|
|
/// @details A master covers the thread that installs it. `derive_party(i)`
|
||
|
|
/// gives party `i` of a run its own stream, keyed by SHA-256 of the
|
||
|
|
/// master and `i`, and `app::run_parties` installs one on each party
|
||
|
|
/// thread, so replaying a master replays every party. Kernels handed
|
||
|
|
/// to a compute pool draw from their party's stream
|
||
|
|
/// (`dpf/thread_work.hpp`). Every CSV row carries the invocation id
|
||
|
|
/// of the process that wrote it (`log::invocation_id()`), and a file
|
||
|
|
/// whose header does not match this layout is moved aside rather
|
||
|
|
/// than appended to.
|
||
|
|
#ifndef LIBDPF_INCLUDE_DPF_EXPERIMENT_HPP__
|
||
|
|
#define LIBDPF_INCLUDE_DPF_EXPERIMENT_HPP__
|
||
|
|
|
||
|
|
#include <algorithm>
|
||
|
|
#include <array>
|
||
|
|
#include <chrono>
|
||
|
|
#include <cstddef>
|
||
|
|
#include <cstdint>
|
||
|
|
#include <cstdio>
|
||
|
|
#include <cstdlib>
|
||
|
|
#include <cstring>
|
||
|
|
#include <filesystem>
|
||
|
|
#include <fstream>
|
||
|
|
#include <iomanip>
|
||
|
|
#include <map>
|
||
|
|
#include <sstream>
|
||
|
|
#include <stdexcept>
|
||
|
|
#include <string>
|
||
|
|
#include <system_error>
|
||
|
|
#include <utility>
|
||
|
|
#include <vector>
|
||
|
|
|
||
|
|
#include <time.h>
|
||
|
|
|
||
|
|
#include "dpf/compose.hpp"
|
||
|
|
#include "dpf/experiment_note.hpp"
|
||
|
|
#include "dpf/log.hpp"
|
||
|
|
#include "dpf/prg_aes.hpp"
|
||
|
|
#include "dpf/prg_count.hpp"
|
||
|
|
#include "dpf/protocol.hpp"
|
||
|
|
#include "dpf/random.hpp"
|
||
|
|
#include "dpf/thread_work.hpp"
|
||
|
|
|
||
|
|
namespace dpf
|
||
|
|
{
|
||
|
|
|
||
|
|
using protocol::round_event;
|
||
|
|
using protocol::round_probe;
|
||
|
|
|
||
|
|
/// @brief Per-thread experiment: master seed stream + cost meter + CSV export.
|
||
|
|
class experiment : public detail::experiment_seed_sink
|
||
|
|
{
|
||
|
|
public:
|
||
|
|
static constexpr std::size_t master_bytes = 32;
|
||
|
|
|
||
|
|
using master_seed = std::array<std::uint8_t, master_bytes>;
|
||
|
|
|
||
|
|
/// @brief Where a master came from.
|
||
|
|
enum class origin : unsigned char
|
||
|
|
{
|
||
|
|
fresh, ///< drawn from OS entropy
|
||
|
|
provided, ///< given to `replay`
|
||
|
|
derived ///< `derive_party` of another master
|
||
|
|
};
|
||
|
|
|
||
|
|
struct noted_seed
|
||
|
|
{
|
||
|
|
std::string name;
|
||
|
|
std::vector<std::uint8_t> bytes;
|
||
|
|
};
|
||
|
|
|
||
|
|
/// @brief Fresh master seed from the system entropy source.
|
||
|
|
/// @details Always reads the OS entropy path (never an outer experiment
|
||
|
|
/// hook), so nested contexts do not steal the parent stream.
|
||
|
|
explicit experiment(std::string name, std::string party = "p0")
|
||
|
|
: name_(std::move(name)), party_(std::move(party))
|
||
|
|
{
|
||
|
|
draw_system_master_();
|
||
|
|
// Master itself must not count as protocol random consumption.
|
||
|
|
reset_random_bytes_count();
|
||
|
|
prg::reset_eval_count();
|
||
|
|
install_();
|
||
|
|
note("master", master_.data(), master_.size());
|
||
|
|
log_master_();
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Replay a previously recorded master seed.
|
||
|
|
static experiment replay(std::string name, const master_seed & seed,
|
||
|
|
std::string party = "p0")
|
||
|
|
{
|
||
|
|
return experiment(std::move(name), std::move(party), seed, origin::provided);
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief The stream party `party` of a run draws from, installed on the
|
||
|
|
/// calling thread (call it on that party's thread). Its master is
|
||
|
|
/// SHA-256 of this master and the party index, so one master
|
||
|
|
/// always yields the same party streams.
|
||
|
|
experiment derive_party(unsigned party) const
|
||
|
|
{
|
||
|
|
std::uint8_t in[master_bytes + 4];
|
||
|
|
std::memcpy(in, master_.data(), master_bytes);
|
||
|
|
for (int i = 0; i < 4; ++i)
|
||
|
|
in[master_bytes + static_cast<std::size_t>(i)] =
|
||
|
|
static_cast<std::uint8_t>((party >> (8 * i)) & 0xffu);
|
||
|
|
const auto digest = log::detail::sha256("libdpf experiment party", in, sizeof(in));
|
||
|
|
master_seed derived{};
|
||
|
|
std::memcpy(derived.data(), digest.data(), derived.size());
|
||
|
|
return experiment(name_, "p" + std::to_string(party), derived, origin::derived);
|
||
|
|
}
|
||
|
|
|
||
|
|
experiment(const experiment &) = delete;
|
||
|
|
experiment & operator=(const experiment &) = delete;
|
||
|
|
|
||
|
|
experiment(experiment && other) noexcept
|
||
|
|
: name_(std::move(other.name_)),
|
||
|
|
party_(std::move(other.party_)),
|
||
|
|
run_id_(other.run_id_),
|
||
|
|
master_(other.master_),
|
||
|
|
origin_(other.origin_),
|
||
|
|
aes_seed_(other.aes_seed_),
|
||
|
|
ctr_(other.ctr_),
|
||
|
|
buf_pos_(other.buf_pos_),
|
||
|
|
prev_hook_(other.prev_hook_),
|
||
|
|
prev_ctx_(other.prev_ctx_),
|
||
|
|
prev_sink_(other.prev_sink_),
|
||
|
|
installed_(other.installed_),
|
||
|
|
seeds_(std::move(other.seeds_)),
|
||
|
|
rounds_(std::move(other.rounds_)),
|
||
|
|
edge_in_(std::move(other.edge_in_)),
|
||
|
|
edge_out_(std::move(other.edge_out_)),
|
||
|
|
edge_plan_out_(std::move(other.edge_plan_out_)),
|
||
|
|
interactive_rounds_(other.interactive_rounds_),
|
||
|
|
dag_depth_(other.dag_depth_),
|
||
|
|
critical_path_(std::move(other.critical_path_)),
|
||
|
|
wall_ns_(other.wall_ns_),
|
||
|
|
cpu_ns_(other.cpu_ns_),
|
||
|
|
prg_evals_(other.prg_evals_),
|
||
|
|
sym_(other.sym_),
|
||
|
|
random_bytes_(other.random_bytes_),
|
||
|
|
bytes_in_(other.bytes_in_),
|
||
|
|
bytes_out_(other.bytes_out_),
|
||
|
|
plan_bytes_out_(other.plan_bytes_out_),
|
||
|
|
config_(std::move(other.config_)),
|
||
|
|
trials_(std::move(other.trials_)),
|
||
|
|
party_trials_(std::move(other.party_trials_)),
|
||
|
|
wire_(other.wire_),
|
||
|
|
timing_started_(other.timing_started_),
|
||
|
|
wall0_(other.wall0_),
|
||
|
|
cpu0_(other.cpu0_)
|
||
|
|
{
|
||
|
|
std::memcpy(buf_, other.buf_, sizeof(buf_));
|
||
|
|
other.installed_ = false;
|
||
|
|
if (installed_)
|
||
|
|
{
|
||
|
|
detail::uniform_bytes_hook = &experiment::hook_;
|
||
|
|
detail::uniform_bytes_ctx = this;
|
||
|
|
detail::experiment_seed_sink_tls = this;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
~experiment() override { uninstall_(); }
|
||
|
|
|
||
|
|
const std::string & name() const noexcept { return name_; }
|
||
|
|
const std::string & party() const noexcept { return party_; }
|
||
|
|
const master_seed & seed() const noexcept { return master_; }
|
||
|
|
std::uint64_t run_id() const noexcept { return run_id_; }
|
||
|
|
|
||
|
|
void set_run_id(std::uint64_t id) noexcept { run_id_ = id; }
|
||
|
|
void set_party(std::string party) { party_ = std::move(party); }
|
||
|
|
|
||
|
|
/// @brief Public alias of `note` for call-site labels.
|
||
|
|
void note_seed(const char * name, const void * bytes, std::size_t n)
|
||
|
|
{
|
||
|
|
note(name, static_cast<const std::uint8_t *>(bytes), n);
|
||
|
|
}
|
||
|
|
|
||
|
|
template <typename T>
|
||
|
|
void note_seed(const char * name, const T & seed)
|
||
|
|
{
|
||
|
|
note_seed(name, &seed, sizeof(seed));
|
||
|
|
}
|
||
|
|
|
||
|
|
void note(const char * name, const std::uint8_t * bytes,
|
||
|
|
std::size_t n) override
|
||
|
|
{
|
||
|
|
if (name == nullptr || bytes == nullptr || n == 0)
|
||
|
|
return;
|
||
|
|
seeds_.push_back({name, std::vector<std::uint8_t>(bytes, bytes + n)});
|
||
|
|
if (std::strcmp(name, "master") != 0)
|
||
|
|
DPF_LOG(debug, "seed").kv("name", name).kv("experiment", name_)
|
||
|
|
.kv("party", party_).kv("source", "noted").kv("bytes", n)
|
||
|
|
.seed("value", bytes, n);
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Noted seeds in order, starting with `master`.
|
||
|
|
const std::vector<noted_seed> & seeds() const noexcept { return seeds_; }
|
||
|
|
|
||
|
|
/// @brief Append a party stream's noted seeds, each name prefixed by
|
||
|
|
/// `party` (`p1/master`, `p1/buffered_prg`).
|
||
|
|
void fold_seeds(const std::string & party, const std::vector<noted_seed> & seeds)
|
||
|
|
{
|
||
|
|
for (const auto & s : seeds)
|
||
|
|
seeds_.push_back({party + "/" + s.name, s.bytes});
|
||
|
|
}
|
||
|
|
|
||
|
|
origin seed_origin() const noexcept { return origin_; }
|
||
|
|
|
||
|
|
/// @brief True when the master was not drawn fresh (`replay` or
|
||
|
|
/// `derive_party`).
|
||
|
|
bool seed_provided() const noexcept { return origin_ != origin::fresh; }
|
||
|
|
|
||
|
|
static const char * origin_name(origin o) noexcept
|
||
|
|
{
|
||
|
|
switch (o)
|
||
|
|
{
|
||
|
|
case origin::fresh:
|
||
|
|
return "fresh";
|
||
|
|
case origin::provided:
|
||
|
|
return "provided";
|
||
|
|
case origin::derived:
|
||
|
|
return "derived";
|
||
|
|
}
|
||
|
|
return "fresh";
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Symmetric-key blocks counted between `begin_timing` and
|
||
|
|
/// `end_timing`, by purpose and primitive (see `prg_count.hpp`).
|
||
|
|
const prg::counts & sym_counts() const noexcept { return sym_; }
|
||
|
|
|
||
|
|
/// @brief Record static plan facts (chain length, per-edge schedule bytes).
|
||
|
|
void ingest_plan(const protocol::plan & p)
|
||
|
|
{
|
||
|
|
interactive_rounds_ = p.rounds();
|
||
|
|
dag_depth_ = p.waves();
|
||
|
|
plan_bytes_out_ = 0;
|
||
|
|
for (auto n : p.slot_bytes_all())
|
||
|
|
plan_bytes_out_ += n;
|
||
|
|
edge_plan_out_.clear();
|
||
|
|
for (std::size_t wi = 0; wi < p.waves(); ++wi)
|
||
|
|
{
|
||
|
|
const auto & w = p.wave(wi);
|
||
|
|
if (w.exchanges.empty())
|
||
|
|
continue;
|
||
|
|
const auto ch = protocol::detail::wave_channel(p, w);
|
||
|
|
edge_plan_out_[static_cast<int>(ch)] += w.slot_bytes;
|
||
|
|
}
|
||
|
|
critical_path_ = build_critical_path_(p);
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Probe that records live round events into this experiment.
|
||
|
|
round_probe probe() noexcept
|
||
|
|
{
|
||
|
|
return round_probe{this, &experiment::on_round_};
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Start wall/CPU/PRG/random timers (call before drive).
|
||
|
|
void begin_timing()
|
||
|
|
{
|
||
|
|
prg::reset_eval_count();
|
||
|
|
reset_random_bytes_count();
|
||
|
|
wall0_ = steady_now_();
|
||
|
|
cpu0_ = thread_cpu_now_();
|
||
|
|
timing_started_ = true;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Stop timers and fold totals.
|
||
|
|
void end_timing()
|
||
|
|
{
|
||
|
|
if (!timing_started_)
|
||
|
|
return;
|
||
|
|
wall_ns_ += steady_now_() - wall0_;
|
||
|
|
cpu_ns_ += thread_cpu_now_() - cpu0_;
|
||
|
|
prg_evals_ += prg::eval_count();
|
||
|
|
const auto sym = prg::snapshot();
|
||
|
|
for (std::size_t i = 0; i < sym.size(); ++i)
|
||
|
|
sym_[i] += sym[i];
|
||
|
|
random_bytes_ += random_bytes_count();
|
||
|
|
timing_started_ = false;
|
||
|
|
}
|
||
|
|
|
||
|
|
std::uint64_t wall_ns() const noexcept { return wall_ns_; }
|
||
|
|
std::uint64_t cpu_ns() const noexcept { return cpu_ns_; }
|
||
|
|
std::uint64_t prg_evals() const noexcept { return prg_evals_; }
|
||
|
|
std::uint64_t random_bytes() const noexcept { return random_bytes_; }
|
||
|
|
std::size_t bytes_in() const noexcept { return bytes_in_; }
|
||
|
|
std::size_t bytes_out() const noexcept { return bytes_out_; }
|
||
|
|
std::size_t interactive_rounds() const noexcept
|
||
|
|
{
|
||
|
|
return interactive_rounds_;
|
||
|
|
}
|
||
|
|
std::size_t dag_depth() const noexcept { return dag_depth_; }
|
||
|
|
std::size_t plan_bytes_out() const noexcept { return plan_bytes_out_; }
|
||
|
|
const std::vector<round_event> & rounds() const noexcept { return rounds_; }
|
||
|
|
std::size_t seed_count() const noexcept { return seeds_.size(); }
|
||
|
|
std::size_t critical_path_length() const noexcept
|
||
|
|
{
|
||
|
|
return critical_path_.size();
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief True if a seed named `name` was noted (including `"master"`).
|
||
|
|
bool has_seed_named(const char * name) const
|
||
|
|
{
|
||
|
|
if (name == nullptr)
|
||
|
|
return false;
|
||
|
|
for (const auto & s : seeds_)
|
||
|
|
if (s.name == name)
|
||
|
|
return true;
|
||
|
|
return false;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Schedule bytes attributed to `channel` by `ingest_plan`.
|
||
|
|
std::size_t plan_edge_bytes(protocol::edge_channel channel) const
|
||
|
|
{
|
||
|
|
return edge_get_(edge_plan_out_, static_cast<int>(channel));
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Live bytes in/out attributed to `channel` by the round probe.
|
||
|
|
std::size_t edge_bytes_in(protocol::edge_channel channel) const
|
||
|
|
{
|
||
|
|
return edge_get_(edge_in_, static_cast<int>(channel));
|
||
|
|
}
|
||
|
|
std::size_t edge_bytes_out(protocol::edge_channel channel) const
|
||
|
|
{
|
||
|
|
return edge_get_(edge_out_, static_cast<int>(channel));
|
||
|
|
}
|
||
|
|
|
||
|
|
std::string seed_hex() const { return to_hex_(master_.data(), master_.size()); }
|
||
|
|
|
||
|
|
/// @brief On-the-wire counters for party 0's links (headers included).
|
||
|
|
struct wire_counts
|
||
|
|
{
|
||
|
|
std::uint64_t bytes_out = 0;
|
||
|
|
std::uint64_t bytes_in = 0;
|
||
|
|
std::uint64_t payload_out = 0;
|
||
|
|
std::uint64_t payload_in = 0;
|
||
|
|
std::uint64_t frames_out = 0;
|
||
|
|
std::uint64_t frames_in = 0;
|
||
|
|
std::uint64_t write_calls = 0;
|
||
|
|
};
|
||
|
|
|
||
|
|
/// @brief Run configuration recorded in `config.csv` (key, value).
|
||
|
|
void set_config(std::vector<std::pair<std::string, std::string>> kv)
|
||
|
|
{
|
||
|
|
config_ = std::move(kv);
|
||
|
|
}
|
||
|
|
const std::vector<std::pair<std::string, std::string>> & config() const noexcept
|
||
|
|
{
|
||
|
|
return config_;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief One timed trial: party 0's wall time and, when given, every
|
||
|
|
/// party's (`party_walls[i]` is party `i`). Recorded in `trials.csv`.
|
||
|
|
void add_trial(std::uint64_t wall_ns,
|
||
|
|
const std::vector<std::uint64_t> & party_walls = {})
|
||
|
|
{
|
||
|
|
trials_.push_back(wall_ns);
|
||
|
|
party_trials_.push_back(party_walls);
|
||
|
|
}
|
||
|
|
const std::vector<std::uint64_t> & trials() const noexcept { return trials_; }
|
||
|
|
|
||
|
|
/// @brief Median of party 0's trial wall times (0 when none).
|
||
|
|
std::uint64_t median_trial_ns() const { return median_(trials_); }
|
||
|
|
|
||
|
|
/// @brief Median over trials of the slowest party's wall time (party 0's
|
||
|
|
/// when a trial did not record the others).
|
||
|
|
std::uint64_t slowest_median_ns() const
|
||
|
|
{
|
||
|
|
std::vector<std::uint64_t> slowest;
|
||
|
|
for (std::size_t t = 0; t < trials_.size(); ++t)
|
||
|
|
{
|
||
|
|
std::uint64_t w = trials_[t];
|
||
|
|
if (t < party_trials_.size())
|
||
|
|
for (auto p : party_trials_[t])
|
||
|
|
w = std::max(w, p);
|
||
|
|
slowest.push_back(w);
|
||
|
|
}
|
||
|
|
return median_(slowest);
|
||
|
|
}
|
||
|
|
|
||
|
|
void set_wire(const wire_counts & w) { wire_ = w; }
|
||
|
|
const wire_counts & wire() const noexcept { return wire_; }
|
||
|
|
|
||
|
|
/// @brief Write / append CSV tables under `dir` (created if missing).
|
||
|
|
void write_csv(const std::string & dir) const
|
||
|
|
{
|
||
|
|
if (dir.empty())
|
||
|
|
throw std::invalid_argument("experiment::write_csv empty dir");
|
||
|
|
std::error_code ec;
|
||
|
|
std::filesystem::create_directories(dir, ec);
|
||
|
|
if (ec)
|
||
|
|
throw std::runtime_error("experiment: cannot create '" + dir + "': "
|
||
|
|
+ ec.message());
|
||
|
|
write_summary_(dir);
|
||
|
|
write_rounds_(dir);
|
||
|
|
write_edges_(dir);
|
||
|
|
write_seeds_(dir);
|
||
|
|
write_critical_path_(dir);
|
||
|
|
write_config_(dir);
|
||
|
|
write_trials_(dir);
|
||
|
|
write_wire_(dir);
|
||
|
|
write_sym_(dir);
|
||
|
|
write_runs_(dir);
|
||
|
|
DPF_LOG(info, "csv").kv("dir", dir).kv("experiment", name_)
|
||
|
|
.kv("party", party_).kv("run_id", run_id_).kv("rounds", rounds_.size())
|
||
|
|
.kv("trials", trials_.size()).kv("seeds", seeds_.size());
|
||
|
|
}
|
||
|
|
|
||
|
|
private:
|
||
|
|
experiment(std::string name, std::string party, master_seed seed, origin o)
|
||
|
|
: name_(std::move(name)), party_(std::move(party)), master_(seed), origin_(o)
|
||
|
|
{
|
||
|
|
reset_random_bytes_count();
|
||
|
|
prg::reset_eval_count();
|
||
|
|
install_();
|
||
|
|
note("master", master_.data(), master_.size());
|
||
|
|
log_master_();
|
||
|
|
}
|
||
|
|
|
||
|
|
static const char * entropy_name_() noexcept
|
||
|
|
{
|
||
|
|
#if defined(LIBDPF_USE_ARC4RANDOM)
|
||
|
|
return "arc4random";
|
||
|
|
#elif defined(LIBDPF_USE_DEV_RANDOM)
|
||
|
|
return "/dev/random";
|
||
|
|
#else
|
||
|
|
return "/dev/urandom";
|
||
|
|
#endif
|
||
|
|
}
|
||
|
|
|
||
|
|
void log_master_() const
|
||
|
|
{
|
||
|
|
const bool derived = origin_ == origin::derived;
|
||
|
|
if (!log::enabled(derived ? log::level::debug : log::level::info))
|
||
|
|
return;
|
||
|
|
log::record(derived ? log::level::debug : log::level::info, "seed")
|
||
|
|
.kv("name", "master").kv("experiment", name_).kv("party", party_)
|
||
|
|
.kv("source", origin_name(origin_))
|
||
|
|
.kv("entropy", origin_ == origin::fresh ? entropy_name_()
|
||
|
|
: derived ? "sha256(master,party)"
|
||
|
|
: "replay")
|
||
|
|
.kv("bytes", master_.size()).seed("value", master_.data(), master_.size());
|
||
|
|
}
|
||
|
|
|
||
|
|
void draw_system_master_()
|
||
|
|
{
|
||
|
|
// Bypass any installed hook so the master is true OS entropy.
|
||
|
|
auto * saved = detail::uniform_bytes_hook;
|
||
|
|
detail::uniform_bytes_hook = nullptr;
|
||
|
|
for (auto & b : master_)
|
||
|
|
uniform_fill(b);
|
||
|
|
detail::uniform_bytes_hook = saved;
|
||
|
|
}
|
||
|
|
|
||
|
|
void install_()
|
||
|
|
{
|
||
|
|
prev_hook_ = detail::uniform_bytes_hook;
|
||
|
|
prev_ctx_ = detail::uniform_bytes_ctx;
|
||
|
|
prev_sink_ = detail::experiment_seed_sink_tls;
|
||
|
|
detail::uniform_bytes_hook = &experiment::hook_;
|
||
|
|
detail::uniform_bytes_ctx = this;
|
||
|
|
detail::experiment_seed_sink_tls = this;
|
||
|
|
installed_ = true;
|
||
|
|
// Derive AES key from the first 16 master bytes.
|
||
|
|
std::memcpy(&aes_seed_, master_.data(), sizeof(aes_seed_));
|
||
|
|
ctr_ = 0;
|
||
|
|
buf_pos_ = sizeof(buf_); // force refill
|
||
|
|
}
|
||
|
|
|
||
|
|
void uninstall_()
|
||
|
|
{
|
||
|
|
if (!installed_)
|
||
|
|
return;
|
||
|
|
if (detail::uniform_bytes_ctx == this)
|
||
|
|
detail::uniform_bytes_ctx = prev_ctx_;
|
||
|
|
if (detail::uniform_bytes_hook == &experiment::hook_)
|
||
|
|
detail::uniform_bytes_hook = prev_hook_;
|
||
|
|
if (detail::experiment_seed_sink_tls == this)
|
||
|
|
detail::experiment_seed_sink_tls = prev_sink_;
|
||
|
|
installed_ = false;
|
||
|
|
}
|
||
|
|
|
||
|
|
static void hook_(void * dst, std::size_t n)
|
||
|
|
{
|
||
|
|
auto * self = static_cast<experiment *>(detail::uniform_bytes_ctx);
|
||
|
|
if (self == nullptr)
|
||
|
|
throw std::logic_error("experiment hook without active context");
|
||
|
|
self->fill_(dst, n);
|
||
|
|
}
|
||
|
|
|
||
|
|
void fill_(void * dst, std::size_t n)
|
||
|
|
{
|
||
|
|
auto * out = static_cast<std::uint8_t *>(dst);
|
||
|
|
while (n > 0)
|
||
|
|
{
|
||
|
|
if (buf_pos_ >= sizeof(buf_))
|
||
|
|
{
|
||
|
|
const prg::purpose_scope harness(prg::purpose::harness);
|
||
|
|
auto blk = prg::aes128::eval(aes_seed_,
|
||
|
|
static_cast<psnip_uint32_t>(ctr_++));
|
||
|
|
std::memcpy(buf_, &blk, sizeof(buf_));
|
||
|
|
buf_pos_ = 0;
|
||
|
|
}
|
||
|
|
const std::size_t take =
|
||
|
|
std::min(n, sizeof(buf_) - buf_pos_);
|
||
|
|
std::memcpy(out, buf_ + buf_pos_, take);
|
||
|
|
buf_pos_ += take;
|
||
|
|
out += take;
|
||
|
|
n -= take;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
static void on_round_(void * ctx, const round_event & ev)
|
||
|
|
{
|
||
|
|
static_cast<experiment *>(ctx)->record_round_(ev);
|
||
|
|
}
|
||
|
|
|
||
|
|
void record_round_(const round_event & ev)
|
||
|
|
{
|
||
|
|
rounds_.push_back(ev);
|
||
|
|
bytes_in_ += ev.bytes_in;
|
||
|
|
bytes_out_ += ev.bytes_out;
|
||
|
|
edge_in_[static_cast<int>(ev.channel)] += ev.bytes_in;
|
||
|
|
edge_out_[static_cast<int>(ev.channel)] += ev.bytes_out;
|
||
|
|
// wall / cpu / prg / random totals come from begin_timing/end_timing
|
||
|
|
// so finish_schedule local work is included once.
|
||
|
|
}
|
||
|
|
|
||
|
|
static std::uint64_t steady_now_()
|
||
|
|
{
|
||
|
|
using clock = std::chrono::steady_clock;
|
||
|
|
return static_cast<std::uint64_t>(
|
||
|
|
std::chrono::duration_cast<std::chrono::nanoseconds>(
|
||
|
|
clock::now().time_since_epoch())
|
||
|
|
.count());
|
||
|
|
}
|
||
|
|
|
||
|
|
static std::uint64_t thread_cpu_now_()
|
||
|
|
{
|
||
|
|
if (have_thread_cpu_clock())
|
||
|
|
return thread_cpu_ns();
|
||
|
|
return steady_now_();
|
||
|
|
}
|
||
|
|
|
||
|
|
static std::uint64_t median_(std::vector<std::uint64_t> v)
|
||
|
|
{
|
||
|
|
if (v.empty())
|
||
|
|
return 0;
|
||
|
|
std::sort(v.begin(), v.end());
|
||
|
|
return v[v.size() / 2];
|
||
|
|
}
|
||
|
|
|
||
|
|
static std::string to_hex_(const std::uint8_t * p, std::size_t n)
|
||
|
|
{
|
||
|
|
std::ostringstream os;
|
||
|
|
os << std::hex << std::setfill('0');
|
||
|
|
for (std::size_t i = 0; i < n; ++i)
|
||
|
|
os << std::setw(2) << static_cast<unsigned>(p[i]);
|
||
|
|
return os.str();
|
||
|
|
}
|
||
|
|
|
||
|
|
static const char * channel_name_(protocol::edge_channel c)
|
||
|
|
{
|
||
|
|
switch (c)
|
||
|
|
{
|
||
|
|
case protocol::edge_channel::peer:
|
||
|
|
return "peer";
|
||
|
|
case protocol::edge_channel::rss_next:
|
||
|
|
return "rss_next";
|
||
|
|
case protocol::edge_channel::dealer:
|
||
|
|
return "dealer";
|
||
|
|
}
|
||
|
|
return "edge";
|
||
|
|
}
|
||
|
|
|
||
|
|
struct path_node
|
||
|
|
{
|
||
|
|
std::uint32_t id = 0;
|
||
|
|
std::size_t wave = 0;
|
||
|
|
std::uint32_t opcode = 0;
|
||
|
|
int effect = 0;
|
||
|
|
};
|
||
|
|
|
||
|
|
static std::vector<path_node> build_critical_path_(const protocol::plan & p)
|
||
|
|
{
|
||
|
|
std::vector<path_node> path;
|
||
|
|
if (p.nodes().empty())
|
||
|
|
return path;
|
||
|
|
std::uint32_t tip = p.nodes().front().id;
|
||
|
|
std::size_t best_w = p.wave_of(protocol::node{tip});
|
||
|
|
for (auto n : p.nodes())
|
||
|
|
{
|
||
|
|
const auto w = p.wave_of(n);
|
||
|
|
if (w >= best_w)
|
||
|
|
{
|
||
|
|
best_w = w;
|
||
|
|
tip = n.id;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
for (;;)
|
||
|
|
{
|
||
|
|
path_node pn;
|
||
|
|
pn.id = tip;
|
||
|
|
pn.wave = p.wave_of(protocol::node{tip});
|
||
|
|
pn.opcode = p.opcode_of(tip);
|
||
|
|
pn.effect = static_cast<int>(p.effect_of(tip));
|
||
|
|
path.push_back(pn);
|
||
|
|
const auto & ins = p.inputs_of(tip);
|
||
|
|
if (ins.empty())
|
||
|
|
break;
|
||
|
|
std::uint32_t next = ins.front();
|
||
|
|
std::size_t nw = p.wave_of(protocol::node{next});
|
||
|
|
for (auto in : ins)
|
||
|
|
{
|
||
|
|
const auto w = p.wave_of(protocol::node{in});
|
||
|
|
if (w >= nw)
|
||
|
|
{
|
||
|
|
nw = w;
|
||
|
|
next = in;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if (next == tip)
|
||
|
|
break;
|
||
|
|
tip = next;
|
||
|
|
}
|
||
|
|
std::reverse(path.begin(), path.end());
|
||
|
|
return path;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief Open `dir/file` for append, writing `header` first when the file
|
||
|
|
/// is new. A file with a different header is renamed to
|
||
|
|
/// `<stem>.before-<UTC>.csv` so rows never land under the wrong
|
||
|
|
/// columns.
|
||
|
|
static std::ofstream open_table_(const std::string & dir, const char * file,
|
||
|
|
const char * header)
|
||
|
|
{
|
||
|
|
const std::string path = dir + "/" + file;
|
||
|
|
std::string existing;
|
||
|
|
{
|
||
|
|
std::ifstream in(path);
|
||
|
|
if (in)
|
||
|
|
std::getline(in, existing);
|
||
|
|
}
|
||
|
|
if (!existing.empty() && existing != header)
|
||
|
|
{
|
||
|
|
std::string stamp;
|
||
|
|
for (char c : log::detail::utc_text(std::chrono::system_clock::now()))
|
||
|
|
if (c != '-' && c != ':')
|
||
|
|
stamp += c;
|
||
|
|
const std::string moved = path.substr(0, path.size() - 4) + ".before-"
|
||
|
|
+ stamp + ".csv";
|
||
|
|
if (std::rename(path.c_str(), moved.c_str()) != 0)
|
||
|
|
throw std::runtime_error("experiment: " + path
|
||
|
|
+ " has another column layout and cannot be moved aside");
|
||
|
|
DPF_LOG(warning, "csv.moved").kv("file", path).kv("to", moved)
|
||
|
|
.kv("detail", "its header differs from this build's columns");
|
||
|
|
existing.clear();
|
||
|
|
}
|
||
|
|
std::ofstream out(path, std::ios::app);
|
||
|
|
if (!out)
|
||
|
|
throw std::runtime_error(std::string("experiment: cannot write ") + file);
|
||
|
|
if (existing.empty())
|
||
|
|
out << header << '\n';
|
||
|
|
return out;
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_summary_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "summary.csv",
|
||
|
|
"name,party,run_id,master_seed,interactive_rounds,dag_depth,"
|
||
|
|
"wall_ns,cpu_ns,prg_evals,random_bytes,bytes_in,bytes_out,"
|
||
|
|
"plan_bytes_out,peer_in,peer_out,rss_in,rss_out,dealer_in,"
|
||
|
|
"dealer_out,median_ns,slowest_median_ns,trials,invocation");
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ',' << seed_hex() << ',' << interactive_rounds_ << ','
|
||
|
|
<< dag_depth_ << ',' << wall_ns_ << ',' << cpu_ns_ << ','
|
||
|
|
<< prg_evals_ << ',' << random_bytes_ << ',' << bytes_in_ << ','
|
||
|
|
<< bytes_out_ << ',' << plan_bytes_out_ << ','
|
||
|
|
<< edge_get_(edge_in_, 0) << ',' << edge_get_(edge_out_, 0) << ','
|
||
|
|
<< edge_get_(edge_in_, 1) << ',' << edge_get_(edge_out_, 1) << ','
|
||
|
|
<< edge_get_(edge_in_, 2) << ',' << edge_get_(edge_out_, 2) << ','
|
||
|
|
<< median_trial_ns() << ',' << slowest_median_ns() << ','
|
||
|
|
<< trials_.size() << ',' << log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_rounds_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "rounds.csv",
|
||
|
|
"name,party,run_id,round,edge,channel,bytes_out,bytes_in,"
|
||
|
|
"wall_ns,cpu_ns,prg_evals,random_bytes,invocation");
|
||
|
|
for (const auto & r : rounds_)
|
||
|
|
{
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ',' << r.round << ',' << r.edge << ','
|
||
|
|
<< channel_name_(r.channel) << ',' << r.bytes_out << ','
|
||
|
|
<< r.bytes_in << ',' << r.wall_ns << ',' << r.cpu_ns << ','
|
||
|
|
<< r.prg_evals << ',' << r.random_bytes << ','
|
||
|
|
<< log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_edges_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "edges.csv",
|
||
|
|
"name,party,run_id,channel,bytes_in,bytes_out,plan_bytes_out,invocation");
|
||
|
|
for (int c = 0; c < 3; ++c)
|
||
|
|
{
|
||
|
|
const auto ch = static_cast<protocol::edge_channel>(c);
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ',' << channel_name_(ch) << ','
|
||
|
|
<< edge_get_(edge_in_, c) << ',' << edge_get_(edge_out_, c)
|
||
|
|
<< ',' << edge_get_(edge_plan_out_, c) << ','
|
||
|
|
<< log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_seeds_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "seeds.csv",
|
||
|
|
"name,party,run_id,seed_name,seed_hex,seed_bytes,invocation");
|
||
|
|
for (const auto & s : seeds_)
|
||
|
|
{
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ',' << csv_escape_(s.name) << ','
|
||
|
|
<< to_hex_(s.bytes.data(), s.bytes.size()) << ','
|
||
|
|
<< s.bytes.size() << ',' << log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_critical_path_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "critical_path.csv",
|
||
|
|
"name,party,run_id,step,node_id,wave,opcode,effect,invocation");
|
||
|
|
for (std::size_t i = 0; i < critical_path_.size(); ++i)
|
||
|
|
{
|
||
|
|
const auto & n = critical_path_[i];
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ',' << i << ',' << n.id << ',' << n.wave << ','
|
||
|
|
<< n.opcode << ',' << n.effect << ',' << log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_config_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "config.csv", "name,party,run_id,key,value,invocation");
|
||
|
|
for (const auto & kv : config_)
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ',' << csv_escape_(kv.first) << ','
|
||
|
|
<< csv_escape_(kv.second) << ',' << log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
|
||
|
|
/// @brief One row per party per trial when parties were recorded, one row
|
||
|
|
/// per trial for this experiment's party otherwise.
|
||
|
|
void write_trials_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "trials.csv",
|
||
|
|
"name,party,run_id,trial,wall_ns,invocation");
|
||
|
|
for (std::size_t t = 0; t < trials_.size(); ++t)
|
||
|
|
{
|
||
|
|
const bool all = t < party_trials_.size() && !party_trials_[t].empty();
|
||
|
|
if (!all)
|
||
|
|
{
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ',' << t << ',' << trials_[t] << ','
|
||
|
|
<< log::invocation_id() << '\n';
|
||
|
|
continue;
|
||
|
|
}
|
||
|
|
for (std::size_t p = 0; p < party_trials_[t].size(); ++p)
|
||
|
|
out << csv_escape_(name_) << ",p" << p << ',' << run_id_ << ','
|
||
|
|
<< t << ',' << party_trials_[t][p] << ','
|
||
|
|
<< log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_wire_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "wire.csv",
|
||
|
|
"name,party,run_id,bytes_out,bytes_in,payload_out,payload_in,"
|
||
|
|
"frames_out,frames_in,write_calls,invocation");
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ',' << run_id_
|
||
|
|
<< ',' << wire_.bytes_out << ',' << wire_.bytes_in << ','
|
||
|
|
<< wire_.payload_out << ',' << wire_.payload_in << ','
|
||
|
|
<< wire_.frames_out << ',' << wire_.frames_in << ','
|
||
|
|
<< wire_.write_calls << ',' << log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_sym_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "sym.csv",
|
||
|
|
"name,party,run_id,purpose,primitive,blocks,invocation");
|
||
|
|
for (std::size_t u = 0; u < prg::purpose_count; ++u)
|
||
|
|
for (std::size_t p = 0; p < prg::primitive_count; ++p)
|
||
|
|
{
|
||
|
|
const auto n = sym_[u * prg::primitive_count + p];
|
||
|
|
if (n == 0)
|
||
|
|
continue;
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ','
|
||
|
|
<< run_id_ << ','
|
||
|
|
<< prg::purpose_name(static_cast<prg::purpose>(u)) << ','
|
||
|
|
<< prg::primitive_name(static_cast<prg::primitive>(p)) << ','
|
||
|
|
<< n << ',' << log::invocation_id() << '\n';
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
void write_runs_(const std::string & dir) const
|
||
|
|
{
|
||
|
|
auto out = open_table_(dir, "runs.csv",
|
||
|
|
"name,party,run_id,invocation,seed_source,written_utc");
|
||
|
|
out << csv_escape_(name_) << ',' << csv_escape_(party_) << ',' << run_id_
|
||
|
|
<< ',' << log::invocation_id() << ',' << origin_name(origin_) << ','
|
||
|
|
<< log::detail::utc_text(std::chrono::system_clock::now()) << '\n';
|
||
|
|
}
|
||
|
|
|
||
|
|
static std::string csv_escape_(const std::string & s)
|
||
|
|
{
|
||
|
|
if (s.find_first_of(",\"\n\r") == std::string::npos)
|
||
|
|
return s;
|
||
|
|
std::string out = "\"";
|
||
|
|
for (char c : s)
|
||
|
|
{
|
||
|
|
if (c == '"')
|
||
|
|
out += "\"\"";
|
||
|
|
else
|
||
|
|
out += c;
|
||
|
|
}
|
||
|
|
out += '"';
|
||
|
|
return out;
|
||
|
|
}
|
||
|
|
|
||
|
|
static std::size_t edge_get_(const std::map<int, std::size_t> & m, int k)
|
||
|
|
{
|
||
|
|
auto it = m.find(k);
|
||
|
|
return it == m.end() ? 0 : it->second;
|
||
|
|
}
|
||
|
|
|
||
|
|
std::string name_;
|
||
|
|
std::string party_;
|
||
|
|
std::uint64_t run_id_ = 0;
|
||
|
|
master_seed master_{};
|
||
|
|
origin origin_ = origin::fresh;
|
||
|
|
prg::aes128::block_type aes_seed_{};
|
||
|
|
std::uint32_t ctr_ = 0;
|
||
|
|
std::uint8_t buf_[16]{};
|
||
|
|
std::size_t buf_pos_ = 16;
|
||
|
|
void (*prev_hook_)(void *, std::size_t) = nullptr;
|
||
|
|
void * prev_ctx_ = nullptr;
|
||
|
|
detail::experiment_seed_sink * prev_sink_ = nullptr;
|
||
|
|
bool installed_ = false;
|
||
|
|
std::vector<noted_seed> seeds_;
|
||
|
|
std::vector<round_event> rounds_;
|
||
|
|
std::map<int, std::size_t> edge_in_;
|
||
|
|
std::map<int, std::size_t> edge_out_;
|
||
|
|
std::map<int, std::size_t> edge_plan_out_;
|
||
|
|
std::size_t interactive_rounds_ = 0;
|
||
|
|
std::size_t dag_depth_ = 0;
|
||
|
|
std::vector<path_node> critical_path_;
|
||
|
|
std::uint64_t wall_ns_ = 0;
|
||
|
|
std::uint64_t cpu_ns_ = 0;
|
||
|
|
std::uint64_t prg_evals_ = 0;
|
||
|
|
prg::counts sym_{};
|
||
|
|
std::uint64_t random_bytes_ = 0;
|
||
|
|
std::size_t bytes_in_ = 0;
|
||
|
|
std::size_t bytes_out_ = 0;
|
||
|
|
std::size_t plan_bytes_out_ = 0;
|
||
|
|
std::vector<std::pair<std::string, std::string>> config_;
|
||
|
|
std::vector<std::uint64_t> trials_;
|
||
|
|
std::vector<std::vector<std::uint64_t>> party_trials_;
|
||
|
|
wire_counts wire_{};
|
||
|
|
bool timing_started_ = false;
|
||
|
|
std::uint64_t wall0_ = 0;
|
||
|
|
std::uint64_t cpu0_ = 0;
|
||
|
|
};
|
||
|
|
|
||
|
|
} // namespace dpf
|
||
|
|
|
||
|
|
#endif // LIBDPF_INCLUDE_DPF_EXPERIMENT_HPP__
|