11#include <benchmark/benchmark.h>
25#include <Monitoring/Monitoring.h>
26#include <fairmq/TransportFactory.h>
55 int quickUpdateInterval = 1;
56 std::vector<MetricSpec> specs{
57 MetricSpec{.
name =
"malformed_inputs", .metricId =
static_cast<short>(ProcessingStatsId::MALFORMED_INPUTS), .minPublishInterval = quickUpdateInterval},
58 MetricSpec{.name =
"dropped_computations", .metricId =
static_cast<short>(ProcessingStatsId::DROPPED_COMPUTATIONS), .minPublishInterval = quickUpdateInterval},
59 MetricSpec{.name =
"dropped_incoming_messages", .metricId =
static_cast<short>(ProcessingStatsId::DROPPED_INCOMING_MESSAGES), .minPublishInterval = quickUpdateInterval},
60 MetricSpec{.name =
"relayed_messages", .metricId =
static_cast<short>(ProcessingStatsId::RELAYED_MESSAGES), .minPublishInterval = quickUpdateInterval}};
61 for (
auto& spec : specs) {
67 r.registerService(ServiceRegistryHelpers::handleForService<DataProcessingStates>(&
states));
68 r.registerService(ServiceRegistryHelpers::handleForService<DataProcessingStats>(&
stats));
69 r.registerService(ServiceRegistryHelpers::handleForService<DriverConfig const>(&
driverConfig));
70 r.registerService(ServiceRegistryHelpers::handleForService<DeviceState>(&
deviceState));
79static void BM_RelayMessageCreation(benchmark::State& state)
88 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
90 for (
auto _ :
state) {
91 fair::mq::MessagePtr header = transport->CreateMessage(
stack.size());
92 fair::mq::MessagePtr payload = transport->CreateMessage(1000);
93 memcpy(header->GetData(),
stack.data(),
stack.size());
101static void BM_RelaySingleSlot(benchmark::State& state)
104 InputSpec spec{
"clusters",
"TPC",
"CLUSTERS"};
106 std::vector<InputRoute> inputs = {
109 std::vector<ForwardRoute> forwards;
110 std::vector<InputChannelInfo> infos{1};
114 relayer.setPipelineLength(4);
125 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
129 std::vector<fair::mq::MessagePtr> inflightMessages;
130 inflightMessages.emplace_back(transport->CreateMessage(
stack.size()));
131 inflightMessages.emplace_back(transport->CreateMessage(1000));
132 memcpy(inflightMessages[0]->GetData(),
stack.data(),
stack.size());
135 for (
auto _ :
state) {
136 relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size());
137 std::vector<RecordAction> ready;
138 relayer.getReadyToProcess(ready);
139 assert(ready.size() == 1);
140 assert(ready[0].slot.index == 0);
141 assert(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
142 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
145 inflightMessages.assign(std::make_move_iterator(
result[0].
begin()),
146 std::make_move_iterator(
result[0].
end()));
153static void BM_RelayMultipleSlots(benchmark::State& state)
156 InputSpec spec{
"clusters",
"TPC",
"CLUSTERS"};
158 std::vector<InputRoute> inputs = {
161 std::vector<ForwardRoute> forwards;
162 std::vector<InputChannelInfo> infos{1};
167 relayer.setPipelineLength(4);
176 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
177 size_t timeslice = 0;
180 Stack placeholder{dh, dph};
185 std::vector<fair::mq::MessagePtr> inflightMessages;
186 inflightMessages.emplace_back(transport->CreateMessage(placeholder.size()));
187 inflightMessages.emplace_back(transport->CreateMessage(1000));
189 for (
auto _ :
state) {
191 memcpy(inflightMessages[0]->GetData(),
stack.data(),
stack.size());
194 relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size());
195 std::vector<RecordAction> ready;
196 relayer.getReadyToProcess(ready);
197 assert(ready.size() == 1);
198 assert(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
199 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
202 inflightMessages.assign(std::make_move_iterator(
result[0].
begin()),
203 std::make_move_iterator(
result[0].
end()));
210static void BM_RelayMultipleRoutes(benchmark::State& state)
213 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
214 InputSpec spec2{
"tracks",
"TPC",
"TRACKS"};
216 std::vector<InputRoute> inputs = {
220 std::vector<ForwardRoute> forwards;
221 std::vector<InputChannelInfo> infos{1};
226 relayer.setPipelineLength(4);
240 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
241 size_t timeslice = 0;
244 Stack stack1{dh1, dph1};
246 std::vector<fair::mq::MessagePtr> inflightMessages;
247 inflightMessages.emplace_back(transport->CreateMessage(stack1.size()));
248 inflightMessages.emplace_back(transport->CreateMessage(1000));
250 memcpy(inflightMessages[0]->GetData(), stack1.data(), stack1.size());
253 Stack stack2{dh2, dph2};
255 inflightMessages.emplace_back(transport->CreateMessage(stack2.size()));
256 inflightMessages.emplace_back(transport->CreateMessage(1000));
258 memcpy(inflightMessages[2]->GetData(), stack2.data(), stack2.size());
260 for (
auto _ :
state) {
262 relayer.relay(inflightMessages[0]->GetData(), &inflightMessages[0], fakeInfo, 2);
263 std::vector<RecordAction> ready;
264 relayer.getReadyToProcess(ready);
265 assert(ready.size() == 1);
266 assert(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
269 relayer.relay(inflightMessages[2]->GetData(), &inflightMessages[2], fakeInfo2, 2);
271 relayer.getReadyToProcess(ready);
272 assert(ready.size() == 1);
273 assert(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
274 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
278 inflightMessages.assign(std::make_move_iterator(
result[0].
begin()),
279 std::make_move_iterator(
result[0].
end()));
280 inflightMessages.insert(inflightMessages.end(),
281 std::make_move_iterator(
result[1].begin()),
282 std::make_move_iterator(
result[1].end()));
289static void BM_RelaySplitParts(benchmark::State& state)
292 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
294 std::vector<InputRoute> inputs = {
298 std::vector<ForwardRoute> forwards;
299 std::vector<InputChannelInfo> infos{1};
304 relayer.setPipelineLength(4);
314 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
315 size_t timeslice = 0;
316 const int nSplitParts =
state.range(0);
318 std::vector<std::unique_ptr<fair::mq::Message>> inflightMessages;
319 inflightMessages.reserve(2 * nSplitParts);
321 for (
size_t i = 0;
i < nSplitParts; ++
i) {
327 fair::mq::MessagePtr header = transport->CreateMessage(
stack.size());
328 fair::mq::MessagePtr payload = transport->CreateMessage(dh.
payloadSize);
330 memcpy(header->GetData(),
stack.data(),
stack.size());
331 inflightMessages.emplace_back(std::move(header));
332 inflightMessages.emplace_back(std::move(payload));
336 for (
auto _ :
state) {
337 relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size());
338 std::vector<RecordAction> ready;
339 relayer.getReadyToProcess(ready);
340 assert(ready.size() == 1);
341 assert(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
342 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
343 inflightMessages.assign(std::make_move_iterator(
result[0].
begin()),
344 std::make_move_iterator(
result[0].
end()));
348BENCHMARK(BM_RelaySplitParts)->Arg(10)->Arg(100)->Arg(1000);
350static void BM_RelayMultiplePayloads(benchmark::State& state)
353 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
355 std::vector<InputRoute> inputs = {
359 std::vector<ForwardRoute> forwards;
360 std::vector<InputChannelInfo> infos{1};
374 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
375 size_t timeslice = 0;
376 const int nPayloads =
state.range(0);
377 std::vector<std::unique_ptr<fair::mq::Message>> inflightMessages;
378 inflightMessages.reserve(nPayloads + 1);
384 fair::mq::MessagePtr header = transport->CreateMessage(
stack.size());
385 memcpy(header->GetData(),
stack.data(),
stack.size());
386 inflightMessages.emplace_back(std::move(header));
387 for (
size_t i = 0;
i < nPayloads; ++
i) {
388 inflightMessages.emplace_back(transport->CreateMessage(dh.
payloadSize));
392 for (
auto _ :
state) {
393 relayer.
relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size(), nPayloads);
394 std::vector<RecordAction> ready;
396 assert(ready.size() == 1);
397 assert(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
399 inflightMessages.assign(std::make_move_iterator(
result[0].begin()),
400 std::make_move_iterator(
result[0].
end()));
404BENCHMARK(BM_RelayMultiplePayloads)->Arg(10)->Arg(100)->Arg(1000);
414static void BM_RelayManyInputs(benchmark::State& state)
417 size_t const nInputs =
state.range(0);
419 std::vector<InputSpec> specs;
420 std::vector<InputRoute> inputs;
421 std::vector<DataHeader> prototypes;
422 specs.reserve(nInputs);
423 for (
size_t i = 0;
i < nInputs; ++
i) {
428 specs.emplace_back(
InputSpec{
"in",
"TST", desc});
436 prototypes.push_back(dh);
438 for (
size_t i = 0;
i < nInputs; ++
i) {
439 inputs.emplace_back(
InputRoute{specs[
i],
i,
"Fake", 0});
442 std::vector<InputChannelInfo> infos{1};
444 auto ref = services.
ref();
445 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
451 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
454 std::vector<fair::mq::MessagePtr> inflight;
455 inflight.reserve(2 * nInputs);
456 for (
size_t i = 0;
i < nInputs; ++
i) {
458 fair::mq::MessagePtr header = transport->CreateMessage(
stack.size());
459 memcpy(header->GetData(),
stack.data(),
stack.size());
460 inflight.emplace_back(std::move(header));
461 inflight.emplace_back(transport->CreateMessage(prototypes[
i].payloadSize));
464 for (
auto _ :
state) {
465 for (
size_t i = 0;
i < nInputs; ++
i) {
467 relayer.
relay(inflight[2 *
i]->GetData(), &inflight[2 *
i], info, 2);
469 std::vector<RecordAction> ready;
471 assert(ready.size() == 1);
474 for (
size_t i = 0;
i < nInputs; ++
i) {
476 inflight.emplace_back(std::move(
msg));
482BENCHMARK(BM_RelayManyInputs)->Arg(1)->Arg(8)->Arg(32)->Arg(128);
header::DataDescription description
o2::monitoring::Monitoring Monitoring
BENCHMARK(BM_RelayMessageCreation)
void getReadyToProcess(std::vector< RecordAction > &completed)
void setPipelineLength(size_t s)
Tune the maximum number of in flight timeslices this can handle.
std::vector< std::vector< fair::mq::MessagePtr > > consumeAllInputsForTimeslice(TimesliceSlot id)
RelayChoice relay(void const *rawHeader, std::unique_ptr< fair::mq::Message > *messages, InputInfo const &info, size_t nMessages, size_t nPayloads=1, OnInsertionCallback onInsertion=nullptr, OnDropCallback onDrop=nullptr)
Defining ITS Vertex explicitly as messageable.
Enum< T >::Iterator begin(Enum< T >)
const DriverConfig driverConfig
DataProcessingStats stats
static constexpr int INVALID
static CompletionPolicy consumeWhenAll(const char *name, CompletionPolicy::Matcher matcher)
Default Completion policy. When all the parts of a record have arrived, consume them.
static CompletionPolicy consumeWhenAny(const char *name, CompletionPolicy::Matcher matcher)
When any of the parts of the record have been received, consume them.
Helper struct to hold statistics about the data processing happening.
void registerMetric(MetricSpec const &spec)
Running state information of a given device.
bool batch
Whether the driver was started in batch mode or not.
void registerService(ServiceTypeHash typeHash, void *service, ServiceKind kind, Salt salt, char const *name=nullptr, ServiceRegistry::SpecIndex specIndex=SpecIndex{-1}) const
static std::function< int64_t(int64_t base, int64_t offset)> defaultCPUTimeConfigurator(uv_loop_t *loop)
static std::function< void(int64_t &base, int64_t &offset)> defaultRealtimeBaseConfigurator(uint64_t offset, uv_loop_t *loop)
uint64_t const void const *restrict const msg