DUNE-DAQ
DUNE Trigger and Data Acquisition software
Toggle main menu visibility
Loading...
Searching...
No Matches
dunedaq
sourcecode
dpdklibs
include
dpdklibs
sender
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
33
#include "
dpdklibs/sender/CallbackGate.hpp
"
34
#include "
dpdklibs/sender/TxEngine.hpp
"
35
#include "
dpdklibs/wibeth/WIBEthPacketBuilder.hpp
"
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
49
namespace
dunedaq::dpdklibs::sender
{
50
51
template
<
typename
Backend>
52
class
WIBEthTransmitter
53
{
54
public
:
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
60
struct
StreamStats
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
74
WIBEthTransmitter
(
const
WIBEthTransmitter
&) =
delete
;
75
WIBEthTransmitter
&
operator=
(
const
WIBEthTransmitter
&) =
delete
;
76
WIBEthTransmitter
(
WIBEthTransmitter
&&) =
delete
;
77
WIBEthTransmitter
&
operator=
(
WIBEthTransmitter
&&) =
delete
;
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;
94
m_backend
=
backend
;
95
}
96
97
// An empty observer disables the pre-transmit byte copy.
98
void
set_post_tx_observer
(
post_tx_observer_t
observer)
99
{
100
std::lock_guard<std::mutex> guard(
m_mutex
);
101
m_post_tx_observer
= std::move(observer);
102
}
103
104
void
set_accepting
(
bool
accepting
) {
m_accepting
.store(
accepting
); }
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
{
120
if
(
m_streams_registered
) {
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
}
134
m_streams_registered
=
true
;
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()) {
144
++
m_unknown_stream_drops
;
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
158
wibeth::HeaderPatch
patch;
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
,
170
wibeth::kEthernetPacketBytes
,
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) {
181
m_post_tx_observer
(
m_observer_scratch
.data(),
wibeth::kEthernetPacketBytes
);
182
}
183
break
;
184
case
TxOutcome::kAllocFailed
:
185
++state.tx_alloc_failures;
186
break
;
187
case
TxOutcome::kPrepareFailed
:
188
++state.tx_prepare_failures;
189
break
;
190
case
TxOutcome::kTxFailed
:
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
208
private
:
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>>();
213
bool
m_streams_registered
=
false
;
214
215
std::atomic<bool>
m_accepting
{
false
};
216
std::atomic<std::uint64_t>
m_unknown_stream_drops
{ 0 };
217
218
Backend
m_backend
{};
219
wibeth::PacketConfig
m_packet_config
{};
220
std::uint64_t
m_det_id
= 0;
221
222
post_tx_observer_t
m_post_tx_observer
;
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_
CallbackGate.hpp
TxEngine.hpp
WIBEthPacketBuilder.hpp
dunedaq::dpdklibs::sender::WIBEthTransmitter::operator=
WIBEthTransmitter & operator=(const WIBEthTransmitter &)=delete
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_streams_registered
bool m_streams_registered
Definition
WIBEthTransmitter.hpp:213
dunedaq::dpdklibs::sender::WIBEthTransmitter::WIBEthTransmitter
WIBEthTransmitter(WIBEthTransmitter &&)=delete
dunedaq::dpdklibs::sender::WIBEthTransmitter::register_streams
void register_streams(Registry ®istry, const std::vector< ConfPtr > &callback_confs)
Definition
WIBEthTransmitter.hpp:118
dunedaq::dpdklibs::sender::WIBEthTransmitter::unknown_stream_drops
std::uint64_t unknown_stream_drops() const
Definition
WIBEthTransmitter.hpp:202
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_mutex
std::mutex m_mutex
Definition
WIBEthTransmitter.hpp:209
dunedaq::dpdklibs::sender::WIBEthTransmitter::backend
Backend & backend()
Definition
WIBEthTransmitter.hpp:206
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_packet_config
wibeth::PacketConfig m_packet_config
Definition
WIBEthTransmitter.hpp:219
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_accepting
std::atomic< bool > m_accepting
Definition
WIBEthTransmitter.hpp:215
dunedaq::dpdklibs::sender::WIBEthTransmitter::tombstone
void tombstone()
Definition
WIBEthTransmitter.hpp:109
dunedaq::dpdklibs::sender::WIBEthTransmitter::operator=
WIBEthTransmitter & operator=(WIBEthTransmitter &&)=delete
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_gate
std::shared_ptr< CallbackGate< WIBEthTransmitter > > m_gate
Definition
WIBEthTransmitter.hpp:212
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_unknown_stream_drops
std::atomic< std::uint64_t > m_unknown_stream_drops
Definition
WIBEthTransmitter.hpp:216
dunedaq::dpdklibs::sender::WIBEthTransmitter::set_accepting
void set_accepting(bool accepting)
Definition
WIBEthTransmitter.hpp:104
dunedaq::dpdklibs::sender::WIBEthTransmitter::post_tx_observer_t
std::function< void(const std::uint8_t *, std::size_t)> post_tx_observer_t
Definition
WIBEthTransmitter.hpp:58
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_streams
std::map< std::string, StreamStats > m_streams
Definition
WIBEthTransmitter.hpp:210
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_backend
Backend m_backend
Definition
WIBEthTransmitter.hpp:218
dunedaq::dpdklibs::sender::WIBEthTransmitter::accepting
bool accepting() const
Definition
WIBEthTransmitter.hpp:105
dunedaq::dpdklibs::sender::WIBEthTransmitter::add_stream
void add_stream(const std::string &callback_id, std::uint64_t initial_seq_id)
Definition
WIBEthTransmitter.hpp:81
dunedaq::dpdklibs::sender::WIBEthTransmitter::configure
void configure(const wibeth::PacketConfig &packet_config, std::uint64_t det_id, const Backend &backend)
Definition
WIBEthTransmitter.hpp:89
dunedaq::dpdklibs::sender::WIBEthTransmitter::set_post_tx_observer
void set_post_tx_observer(post_tx_observer_t observer)
Definition
WIBEthTransmitter.hpp:98
dunedaq::dpdklibs::sender::WIBEthTransmitter::~WIBEthTransmitter
~WIBEthTransmitter()
Definition
WIBEthTransmitter.hpp:72
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_post_tx_observer
post_tx_observer_t m_post_tx_observer
Definition
WIBEthTransmitter.hpp:222
dunedaq::dpdklibs::sender::WIBEthTransmitter::armed
bool armed() const
Definition
WIBEthTransmitter.hpp:110
dunedaq::dpdklibs::sender::WIBEthTransmitter::WIBEthTransmitter
WIBEthTransmitter()=default
dunedaq::dpdklibs::sender::WIBEthTransmitter::snapshot_stats
std::map< std::string, StreamStats > snapshot_stats() const
Definition
WIBEthTransmitter.hpp:196
dunedaq::dpdklibs::sender::WIBEthTransmitter::transmit
void transmit(const std::string &callback_id, PayloadT &&payload)
Definition
WIBEthTransmitter.hpp:139
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_observer_scratch
std::array< std::uint8_t, wibeth::kEthernetPacketBytes > m_observer_scratch
Definition
WIBEthTransmitter.hpp:225
dunedaq::dpdklibs::sender::WIBEthTransmitter::m_det_id
std::uint64_t m_det_id
Definition
WIBEthTransmitter.hpp:220
dunedaq::dpdklibs::sender::WIBEthTransmitter::WIBEthTransmitter
WIBEthTransmitter(const WIBEthTransmitter &)=delete
dunedaq::dpdklibs::sender
Definition
CallbackGate.hpp:35
dunedaq::dpdklibs::sender::send_packet
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
dunedaq::dpdklibs::sender::TxOutcome::kTxFailed
@ kTxFailed
Definition
TxEngine.hpp:45
dunedaq::dpdklibs::sender::TxOutcome::kPrepareFailed
@ kPrepareFailed
Definition
TxEngine.hpp:44
dunedaq::dpdklibs::sender::TxOutcome::kAllocFailed
@ kAllocFailed
Definition
TxEngine.hpp:43
dunedaq::dpdklibs::sender::TxOutcome::kSent
@ kSent
Definition
TxEngine.hpp:42
dunedaq::dpdklibs::wibeth::construct_packet
void construct_packet(void *ethernet_packet, const void *wibeth_frame, const PacketConfig &cfg, const HeaderPatch &patch={})
Definition
WIBEthPacketBuilder.cpp:54
dunedaq::dpdklibs::wibeth::kEthernetPacketBytes
constexpr std::size_t kEthernetPacketBytes
Definition
WIBEthPacketBuilder.hpp:32
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats
Definition
WIBEthTransmitter.hpp:61
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats::tx_alloc_failures
std::uint64_t tx_alloc_failures
Definition
WIBEthTransmitter.hpp:66
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats::sent
std::uint64_t sent
Definition
WIBEthTransmitter.hpp:64
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats::tx_failures
std::uint64_t tx_failures
Definition
WIBEthTransmitter.hpp:68
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats::tx_prepare_failures
std::uint64_t tx_prepare_failures
Definition
WIBEthTransmitter.hpp:67
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats::dropped_not_running
std::uint64_t dropped_not_running
Definition
WIBEthTransmitter.hpp:65
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats::next_seq
std::uint64_t next_seq
Definition
WIBEthTransmitter.hpp:62
dunedaq::dpdklibs::sender::WIBEthTransmitter::StreamStats::accepted
std::uint64_t accepted
Definition
WIBEthTransmitter.hpp:63
dunedaq::dpdklibs::wibeth::HeaderPatch
Definition
WIBEthPacketBuilder.hpp:64
dunedaq::dpdklibs::wibeth::HeaderPatch::set_det_id
bool set_det_id
Definition
WIBEthPacketBuilder.hpp:65
dunedaq::dpdklibs::wibeth::HeaderPatch::force_block_length
bool force_block_length
Definition
WIBEthPacketBuilder.hpp:69
dunedaq::dpdklibs::wibeth::HeaderPatch::seq_id
std::uint64_t seq_id
Definition
WIBEthPacketBuilder.hpp:68
dunedaq::dpdklibs::wibeth::HeaderPatch::set_seq_id
bool set_seq_id
Definition
WIBEthPacketBuilder.hpp:67
dunedaq::dpdklibs::wibeth::HeaderPatch::det_id
std::uint64_t det_id
Definition
WIBEthPacketBuilder.hpp:66
dunedaq::dpdklibs::wibeth::PacketConfig
Definition
WIBEthPacketBuilder.hpp:53
dunedaq::dpdklibs::wibeth::PacketConfig::packet_id
std::uint16_t packet_id
Definition
WIBEthPacketBuilder.hpp:60
Generated on
for DUNE-DAQ by
1.18.0