DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
WIBEthTransmitter.hpp
Go to the documentation of this file.
1
30#ifndef DPDKLIBS_INCLUDE_DPDKLIBS_SENDER_WIBETHTRANSMITTER_HPP_
31#define DPDKLIBS_INCLUDE_DPDKLIBS_SENDER_WIBETHTRANSMITTER_HPP_
32
36
37#include <array>
38#include <atomic>
39#include <cstddef>
40#include <cstdint>
41#include <functional>
42#include <map>
43#include <memory>
44#include <mutex>
45#include <string>
46#include <utility>
47#include <vector>
48
50
51template<typename Backend>
53{
54public:
55 // Called with the transmitted bytes after each successful transmit; packets
56 // the backend rejected are not observed. Runs under the transmitter mutex.
57 // Must not throw. Its execution time adds to the transmit path.
58 using post_tx_observer_t = std::function<void(const std::uint8_t*, std::size_t)>;
59
61 {
62 std::uint64_t next_seq = 0;
63 std::uint64_t accepted = 0;
64 std::uint64_t sent = 0;
65 std::uint64_t dropped_not_running = 0;
66 std::uint64_t tx_alloc_failures = 0;
67 std::uint64_t tx_prepare_failures = 0;
68 std::uint64_t tx_failures = 0;
69 };
70
71 WIBEthTransmitter() = default;
72 ~WIBEthTransmitter() { m_gate->tombstone(); }
73
78
79 // Create per-stream state. Idempotent for an existing id: a repeated call
80 // does not reset live counters.
81 void add_stream(const std::string& callback_id, std::uint64_t initial_seq_id)
82 {
83 std::lock_guard<std::mutex> guard(m_mutex);
84 StreamStats stats;
85 stats.next_seq = initial_seq_id % 4096;
86 m_streams.emplace(callback_id, stats);
87 }
88
89 void configure(const wibeth::PacketConfig& packet_config, std::uint64_t det_id, const Backend& backend)
90 {
91 std::lock_guard<std::mutex> guard(m_mutex);
92 m_packet_config = packet_config;
93 m_det_id = det_id;
95 }
96
97 // An empty observer disables the pre-transmit byte copy.
99 {
100 std::lock_guard<std::mutex> guard(m_mutex);
101 m_post_tx_observer = std::move(observer);
102 }
103
105 bool accepting() const { return m_accepting.load(); }
106
107 // Drains in-flight dispatches, then blocks new ones until the next
108 // register_streams() re-arms the gate.
109 void tombstone() { m_gate->tombstone(); }
110 bool armed() const { return m_gate->armed(); }
111
112 // Registers one gate-dispatching lambda per configuration, then arms the
113 // gate. Registry contract (matched by the pinned
114 // datahandlinglibs::DataMoveCallbackRegistry):
115 // registry.register_callback<PayloadT>(conf, std::function<void(PayloadT&&)>)
116 // keyed by conf->UID().
117 template<typename PayloadT, typename Registry, typename ConfPtr>
118 void register_streams(Registry& registry, const std::vector<ConfPtr>& callback_confs)
119 {
121 m_gate->arm(this);
122 return;
123 }
124 for (const auto& conf : callback_confs) {
125 const std::string callback_id = conf->UID();
126 // The lambda captures the gate, not the transmitter. If the
127 // process-global registry outlives this object, dispatch becomes a drop
128 // rather than a call through a dangling pointer.
129 auto gate = m_gate;
130 registry.template register_callback<PayloadT>(conf, [gate, callback_id](PayloadT&& payload) {
131 gate->dispatch([&](WIBEthTransmitter& owner) { owner.transmit(callback_id, std::move(payload)); });
132 });
133 }
135 m_gate->arm(this);
136 }
137
138 template<typename PayloadT>
139 void transmit(const std::string& callback_id, PayloadT&& payload)
140 {
141 std::lock_guard<std::mutex> guard(m_mutex);
142 auto state_it = m_streams.find(callback_id);
143 if (state_it == m_streams.end()) {
145 return;
146 }
147 auto& state = state_it->second;
148 ++state.accepted;
149
150 if (!m_accepting.load()) {
151 ++state.dropped_not_running;
152 return;
153 }
154
155 auto packet_config = m_packet_config;
156 packet_config.packet_id = static_cast<std::uint16_t>(state.sent & 0xffffU);
157
159 patch.set_det_id = true;
160 patch.det_id = m_det_id;
161 patch.set_seq_id = true;
162 patch.seq_id = state.next_seq;
163 patch.force_block_length = true;
164
165 const bool want_observer = static_cast<bool>(m_post_tx_observer);
166 std::uint8_t* pre_tx_copy = want_observer ? m_observer_scratch.data() : nullptr;
167
168 const auto outcome = send_packet(
169 m_backend,
171 pre_tx_copy,
172 [&](std::uint8_t* dst) { wibeth::construct_packet(dst, payload.data, packet_config, patch); },
173 [&]() {
174 ++state.sent;
175 state.next_seq = (state.next_seq + 1) % 4096;
176 });
177
178 switch (outcome) {
179 case TxOutcome::kSent:
180 if (want_observer) {
182 }
183 break;
185 ++state.tx_alloc_failures;
186 break;
188 ++state.tx_prepare_failures;
189 break;
191 ++state.tx_failures;
192 break;
193 }
194 }
195
196 std::map<std::string, StreamStats> snapshot_stats() const
197 {
198 std::lock_guard<std::mutex> guard(m_mutex);
199 return m_streams;
200 }
201
202 std::uint64_t unknown_stream_drops() const { return m_unknown_stream_drops.load(); }
203
204 // Configuration and inspection access. Not synchronized against in-flight
205 // transmits; call from the owning module's command context or from tests.
206 Backend& backend() { return m_backend; }
207
208private:
209 mutable std::mutex m_mutex;
210 std::map<std::string, StreamStats> m_streams;
211
212 std::shared_ptr<CallbackGate<WIBEthTransmitter>> m_gate = std::make_shared<CallbackGate<WIBEthTransmitter>>();
214
215 std::atomic<bool> m_accepting{ false };
216 std::atomic<std::uint64_t> m_unknown_stream_drops{ 0 };
217
218 Backend m_backend{};
220 std::uint64_t m_det_id = 0;
221
223 // Destination of the pre-transmit byte copy. Guarded by m_mutex, like the
224 // rest of the transmit path.
225 std::array<std::uint8_t, wibeth::kEthernetPacketBytes> m_observer_scratch{};
226};
227
228} // namespace dunedaq::dpdklibs::sender
229
230#endif // DPDKLIBS_INCLUDE_DPDKLIBS_SENDER_WIBETHTRANSMITTER_HPP_
WIBEthTransmitter & operator=(const WIBEthTransmitter &)=delete
WIBEthTransmitter(WIBEthTransmitter &&)=delete
void register_streams(Registry &registry, const std::vector< ConfPtr > &callback_confs)
WIBEthTransmitter & operator=(WIBEthTransmitter &&)=delete
std::shared_ptr< CallbackGate< WIBEthTransmitter > > m_gate
std::atomic< std::uint64_t > m_unknown_stream_drops
std::function< void(const std::uint8_t *, std::size_t)> post_tx_observer_t
std::map< std::string, StreamStats > m_streams
void add_stream(const std::string &callback_id, std::uint64_t initial_seq_id)
void configure(const wibeth::PacketConfig &packet_config, std::uint64_t det_id, const Backend &backend)
void set_post_tx_observer(post_tx_observer_t observer)
std::map< std::string, StreamStats > snapshot_stats() const
void transmit(const std::string &callback_id, PayloadT &&payload)
std::array< std::uint8_t, wibeth::kEthernetPacketBytes > m_observer_scratch
WIBEthTransmitter(const WIBEthTransmitter &)=delete
TxOutcome send_packet(Backend &backend, std::size_t packet_bytes, std::uint8_t *pre_tx_copy, BuildFn &&build, OnSentFn &&on_sent)
Definition TxEngine.hpp:84
void construct_packet(void *ethernet_packet, const void *wibeth_frame, const PacketConfig &cfg, const HeaderPatch &patch={})
constexpr std::size_t kEthernetPacketBytes