libdpf/examples/protocol/stream_app_smoke.cpp

84 lines
2.6 KiB
C++
Raw Permalink Normal View History

/// @file examples/protocol/stream_app_smoke.cpp
/// @brief Thin app-runtime smoke: prep files or memory dealer + compose on streams.
/// @details Drives a composed open via `drive_both_on_streams` (the stream
/// framework). The same schedule runs unchanged over the event-driven
/// backends (`dpf::net::async_stream_array` + `async_round_sink` /
/// `dpf::async::overlapped_byte_protocol`) and, on Linux, over real
/// SCTP (`async_sctp_stream_array`). The `experiment_bench` harness
/// selects the transport with
/// `DPF_TRANSPORT=memory|stream|async|mux|parallel|sctp`.
#include <cstdint>
#include <cstring>
#include <iostream>
#include <vector>
#include "dpf/app_runtime.hpp"
#include "dpf/compose.hpp"
#include "dpf/prep_source.hpp"
#include "dpf/protocol_roles.hpp"
namespace
{
void put_raw(std::vector<std::vector<std::uint8_t>> & values,
dpf::protocol::node n, const void * src, std::size_t nbyte)
{
values[n.id].assign(nbyte, 0);
std::memcpy(values[n.id].data(), src, nbyte);
}
} // namespace
int main()
{
// Prep: file dealer round-trip (roles helper + app_runtime reader).
{
const std::string base = "/tmp/libdpf_stream_app_smoke_prep";
dpf::prep::demand d{};
d.ring_triples = 1;
dpf::roles::dealer_write_files(base, d);
auto c0 = dpf::app::open_file_prep_cursor(base, 0);
auto c1 = dpf::app::open_file_prep_cursor(base, 1);
if (c0.limb() != 8 || c1.limb() != 8)
{
std::cerr << "prep cursor\n";
return 1;
}
}
// Online: memory dealer + single-wave exchange plan on stream arrays.
dpf::protocol::composer c0(0);
dpf::protocol::composer c1(1);
auto x0 = c0.input(dpf::protocol::domain::a, 8);
auto x1 = c1.input(dpf::protocol::domain::a, 8);
auto e0 = c0.exchange(x0);
auto e1 = c1.exchange(x1);
auto p0 = c0.schedule();
auto p1 = c1.schedule();
dpf::prep::demand d{};
d.ring_triples = 1;
auto [d0, d1] = dpf::prep::deal_memory_pair(d);
(void)d0;
(void)d1;
std::vector<std::vector<std::uint8_t>> v0(p0.nodes().size());
std::vector<std::vector<std::uint8_t>> v1(p1.nodes().size());
const std::uint64_t a = 5, b = 9;
put_raw(v0, x0, &a, 8);
put_raw(v1, x1, &b, 8);
dpf::protocol::drive_both_on_streams(p0, p1, v0, v1);
std::uint64_t open0 = 0, open1 = 0;
std::memcpy(&open0, v0[e0.id].data(), 8);
std::memcpy(&open1, v1[e1.id].data(), 8);
if (open0 != 14u || open1 != 14u)
{
std::cerr << "exchange " << open0 << " " << open1 << "\n";
return 1;
}
std::cout << "ok\n";
return 0;
}