97 int quickUpdateInterval = 1;
99 std::vector<MetricSpec> specs{
100 MetricSpec{.
name =
"malformed_inputs", .metricId =
static_cast<short>(ProcessingStatsId::MALFORMED_INPUTS), .minPublishInterval = quickUpdateInterval},
101 MetricSpec{.name =
"dropped_computations", .metricId =
static_cast<short>(ProcessingStatsId::DROPPED_COMPUTATIONS), .minPublishInterval = quickUpdateInterval},
102 MetricSpec{.name =
"dropped_incoming_messages", .metricId =
static_cast<short>(ProcessingStatsId::DROPPED_INCOMING_MESSAGES), .minPublishInterval = quickUpdateInterval},
103 MetricSpec{.name =
"relayed_messages", .metricId =
static_cast<short>(ProcessingStatsId::RELAYED_MESSAGES), .minPublishInterval = quickUpdateInterval}};
105 for (
auto& spec : specs) {
106 stats.registerMetric(spec);
110 ref.registerService(ServiceRegistryHelpers::handleForService<Monitoring>(&monitoring));
111 ref.registerService(ServiceRegistryHelpers::handleForService<DataProcessingStats>(&stats));
112 ref.registerService(ServiceRegistryHelpers::handleForService<DataProcessingStates>(&
states));
113 ref.registerService(ServiceRegistryHelpers::handleForService<DriverConfig const>(&driverConfig));
114 ref.registerService(ServiceRegistryHelpers::handleForService<DeviceState>(&
state));
117 SECTION(
"TestNoWait")
119 InputSpec spec{
"clusters",
"TPC",
"CLUSTERS"};
121 std::vector<InputRoute> inputs = {
124 std::vector<ForwardRoute> forwards;
125 std::vector<InputChannelInfo> infos{1};
127 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
143 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
144 std::array<fair::mq::MessagePtr, 2> messages;
145 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
147 messages[1] = transport->CreateMessage(1000);
148 fair::mq::MessagePtr& header = messages[0];
149 fair::mq::MessagePtr& payload = messages[1];
151 relayer.relay(header->GetData(), messages.data(), fakeInfo, messages.size());
152 std::vector<RecordAction> ready;
153 relayer.getReadyToProcess(ready);
154 REQUIRE(ready.size() == 1);
155 REQUIRE(ready[0].slot.index == 0);
156 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
157 REQUIRE(header.get() ==
nullptr);
158 REQUIRE(payload.get() ==
nullptr);
159 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
166 SECTION(
"TestNoWaitMatcher")
171 std::vector<InputRoute> inputs = {
174 std::vector<ForwardRoute> forwards;
175 std::vector<InputChannelInfo> infos{1};
177 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
193 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
194 std::array<fair::mq::MessagePtr, 2> messages;
195 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
197 messages[1] = transport->CreateMessage(1000);
198 fair::mq::MessagePtr& header = messages[0];
199 fair::mq::MessagePtr& payload = messages[1];
201 relayer.relay(header->GetData(), messages.data(), fakeInfo, messages.size());
202 std::vector<RecordAction> ready;
203 relayer.getReadyToProcess(ready);
204 REQUIRE(ready.size() == 1);
205 REQUIRE(ready[0].slot.index == 0);
206 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
207 REQUIRE(header.get() ==
nullptr);
208 REQUIRE(payload.get() ==
nullptr);
209 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
231 std::vector<InputRoute> inputs = {
235 std::vector<ForwardRoute> forwards;
237 std::vector<InputChannelInfo> infos{1};
239 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
245 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
246 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
248 auto createMessage = [&transport, &channelAlloc, &relayer](
DataHeader& dh,
size_t time) {
249 std::array<fair::mq::MessagePtr, 2> messages;
251 messages[1] = transport->CreateMessage(1000);
252 fair::mq::MessagePtr& header = messages[0];
253 fair::mq::MessagePtr& payload = messages[1];
255 relayer.relay(header->GetData(), messages.data(), fakeInfo, messages.size());
256 REQUIRE(header.get() ==
nullptr);
257 REQUIRE(payload.get() ==
nullptr);
277 createMessage(dh1, 0);
278 std::vector<RecordAction> ready;
279 relayer.getReadyToProcess(ready);
280 REQUIRE(ready.size() == 0);
282 createMessage(dh2, 0);
284 relayer.getReadyToProcess(ready);
285 REQUIRE(ready.size() == 1);
286 REQUIRE(ready[0].slot.index == 0);
287 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
289 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
298 SECTION(
"TestRelayBug")
312 std::vector<InputRoute> inputs = {
316 std::vector<ForwardRoute> forwards;
318 std::vector<InputChannelInfo> infos{1};
320 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
326 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
327 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
329 auto createMessage = [&transport, &channelAlloc, &relayer](
DataHeader& dh,
size_t time) {
330 std::array<fair::mq::MessagePtr, 2> messages;
332 messages[1] = transport->CreateMessage(1000);
333 fair::mq::MessagePtr& header = messages[0];
334 fair::mq::MessagePtr& payload = messages[1];
336 relayer.relay(header->GetData(), messages.data(), fakeInfo, messages.size());
337 REQUIRE(header.get() ==
nullptr);
338 REQUIRE(payload.get() ==
nullptr);
367 createMessage(dh1, 0);
368 std::vector<RecordAction> ready;
369 relayer.getReadyToProcess(ready);
370 REQUIRE(ready.size() == 0);
371 createMessage(dh1, 1);
373 relayer.getReadyToProcess(ready);
374 REQUIRE(ready.size() == 0);
375 createMessage(dh2, 0);
377 relayer.getReadyToProcess(ready);
378 REQUIRE(ready.size() == 1);
379 REQUIRE(ready[0].slot.index == 0);
380 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
381 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
382 createMessage(dh2, 1);
384 relayer.getReadyToProcess(ready);
385 REQUIRE(ready.size() == 1);
386 REQUIRE(ready[0].slot.index == 1);
387 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
388 result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
396 InputSpec spec{
"clusters",
"TPC",
"CLUSTERS"};
398 std::vector<InputRoute> inputs = {
400 std::vector<ForwardRoute> forwards;
403 std::vector<InputChannelInfo> infos{1};
405 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
420 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
421 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
422 auto createMessage = [&transport, &channelAlloc, &relayer, &dh](
auto const&
h) {
423 std::array<fair::mq::MessagePtr, 2> messages;
425 messages[1] = transport->CreateMessage(1000);
426 fair::mq::MessagePtr& header = messages[0];
427 fair::mq::MessagePtr& payload = messages[1];
429 auto res = relayer.relay(header->GetData(), messages.data(), fakeInfo, messages.size());
430 REQUIRE((
res.type != DataRelayer::RelayChoice::Type::WillRelay || header.get() ==
nullptr));
431 REQUIRE((
res.type != DataRelayer::RelayChoice::Type::WillRelay || payload.get() ==
nullptr));
432 REQUIRE((
res.type != DataRelayer::RelayChoice::Type::Backpressured || header.get() !=
nullptr));
433 REQUIRE((
res.type != DataRelayer::RelayChoice::Type::Backpressured || payload.get() !=
nullptr));
439 std::vector<RecordAction> ready;
440 relayer.getReadyToProcess(ready);
441 REQUIRE(ready.size() == 2);
442 REQUIRE(ready[0].slot.index == 1);
443 REQUIRE(ready[1].slot.index == 0);
444 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
445 REQUIRE(ready[1].
op == CompletionPolicy::CompletionOp::Consume);
446 for (
size_t i = 0;
i < ready.size(); ++
i) {
447 auto result = relayer.consumeAllInputsForTimeslice(ready[
i].slot);
455 relayer.getReadyToProcess(ready);
456 REQUIRE(ready.size() == 2);
458 auto result1 = relayer.consumeAllInputsForTimeslice(ready[0].slot);
459 auto result2 = relayer.consumeAllInputsForTimeslice(ready[1].slot);
467 SECTION(
"TestPolicies")
470 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
471 InputSpec spec2{
"tracks",
"TPC",
"TRACKS"};
473 std::vector<InputRoute> inputs = {
478 std::vector<ForwardRoute> forwards;
479 std::vector<InputChannelInfo> infos{1};
481 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
504 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
505 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
506 auto createMessage = [&transport, &channelAlloc, &relayer](
auto const& dh,
auto const&
h) {
507 std::array<fair::mq::MessagePtr, 2> messages;
509 messages[1] = transport->CreateMessage(1000);
510 fair::mq::MessagePtr& header = messages[0];
512 return relayer.relay(header->GetData(), messages.data(), fakeInfo, messages.size());
517 std::vector<RecordAction> ready1;
518 relayer.getReadyToProcess(ready1);
519 REQUIRE(ready1.size() == 1);
520 REQUIRE(ready1[0].slot.index == 0);
521 REQUIRE(ready1[0].
op == CompletionPolicy::CompletionOp::Process);
524 std::vector<RecordAction> ready2;
525 relayer.getReadyToProcess(ready2);
526 REQUIRE(ready2.size() == 1);
527 REQUIRE(ready2[0].slot.index == 1);
528 REQUIRE(ready2[0].
op == CompletionPolicy::CompletionOp::Process);
531 std::vector<RecordAction> ready3;
532 relayer.getReadyToProcess(ready3);
533 REQUIRE(ready3.size() == 1);
534 REQUIRE(ready3[0].slot.index == 1);
535 REQUIRE(ready3[0].
op == CompletionPolicy::CompletionOp::Consume);
542 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
543 InputSpec spec2{
"tracks",
"TPC",
"TRACKS"};
545 std::vector<InputRoute> inputs = {
550 std::vector<ForwardRoute> forwards;
551 std::vector<InputChannelInfo> infos{1};
553 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
576 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
577 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
578 auto createMessage = [&transport, &channelAlloc, &relayer](
auto const& dh,
auto const&
h) {
579 std::array<fair::mq::MessagePtr, 2> messages;
581 messages[1] = transport->CreateMessage(1000);
582 fair::mq::MessagePtr& header = messages[0];
584 return relayer.relay(header->GetData(), messages.data(), fakeInfo, messages.size());
592 std::vector<RecordAction> ready;
593 relayer.getReadyToProcess(ready);
594 REQUIRE(ready.size() == 0);
598 SECTION(
"TestTooMany")
601 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
602 InputSpec spec2{
"tracks",
"TPC",
"TRACKS"};
604 std::vector<InputRoute> inputs = {
609 std::vector<ForwardRoute> forwards;
610 std::vector<InputChannelInfo> infos{1};
612 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
635 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
636 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
638 std::array<fair::mq::MessagePtr, 4> messages;
640 messages[1] = transport->CreateMessage(1000);
641 fair::mq::MessagePtr& header = messages[0];
642 fair::mq::MessagePtr& payload = messages[1];
644 relayer.relay(header->GetData(), &messages[0], fakeInfo, 2);
645 REQUIRE(header.get() ==
nullptr);
646 REQUIRE(payload.get() ==
nullptr);
649 messages[3] = transport->CreateMessage(1000);
650 fair::mq::MessagePtr& header2 = messages[2];
651 fair::mq::MessagePtr& payload2 = messages[3];
653 auto action = relayer.relay(header2->GetData(), &messages[2], fakeInfo2, 2);
654 REQUIRE(action.type == DataRelayer::RelayChoice::Type::Backpressured);
655 REQUIRE(header2.get() !=
nullptr);
656 REQUIRE(payload2.get() !=
nullptr);
659 SECTION(
"SplitParts")
662 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
663 InputSpec spec2{
"its",
"ITS",
"CLUSTERS"};
665 std::vector<InputRoute> inputs = {
670 std::vector<ForwardRoute> forwards;
671 std::vector<InputChannelInfo> infos{1};
673 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
696 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
697 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
699 std::array<fair::mq::MessagePtr, 6> messages;
701 messages[1] = transport->CreateMessage(1000);
702 fair::mq::MessagePtr& header = messages[0];
703 fair::mq::MessagePtr& payload = messages[1];
705 relayer.relay(header->GetData(), &messages[0], fakeInfo, 2);
706 REQUIRE(header.get() ==
nullptr);
707 REQUIRE(payload.get() ==
nullptr);
710 messages[3] = transport->CreateMessage(1000);
711 fair::mq::MessagePtr& header2 = messages[2];
712 fair::mq::MessagePtr& payload2 = messages[3];
714 auto action = relayer.relay(header2->GetData(), &messages[2], fakeInfo, 2);
715 REQUIRE(action.type == DataRelayer::RelayChoice::Type::Backpressured);
716 CHECK(action.timeslice.value == 1);
717 REQUIRE(header2.get() !=
nullptr);
718 REQUIRE(payload2.get() !=
nullptr);
721 messages[5] = transport->CreateMessage(1000);
723 relayer.relay(header2->GetData(), &messages[4], fakeInfo3, 2);
724 REQUIRE(action.type == DataRelayer::RelayChoice::Type::Backpressured);
725 CHECK(action.timeslice.value == 1);
726 REQUIRE(header2.get() !=
nullptr);
727 REQUIRE(payload2.get() !=
nullptr);
730 SECTION(
"SplitPayloadPairs")
733 InputSpec spec1{
"clusters",
"TPC",
"CLUSTERS"};
735 std::vector<InputRoute> inputs = {
739 std::vector<ForwardRoute> forwards;
740 std::vector<InputChannelInfo> infos{1};
742 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
750 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
751 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
752 size_t timeslice = 0;
754 const int nSplitParts = 100;
755 std::vector<std::unique_ptr<fair::mq::Message>> splitParts;
756 splitParts.reserve(2 * nSplitParts);
758 for (
size_t i = 0;
i < nSplitParts; ++
i) {
759 dh.splitPayloadIndex =
i;
760 dh.splitPayloadParts = nSplitParts;
763 fair::mq::MessagePtr payload = transport->CreateMessage(100);
765 splitParts.emplace_back(std::move(header));
766 splitParts.emplace_back(std::move(payload));
768 REQUIRE(splitParts.size() == 2 * nSplitParts);
771 relayer.relay(splitParts[0]->GetData(), splitParts.data(), fakeInfo, splitParts.size());
772 std::vector<RecordAction> ready;
773 relayer.getReadyToProcess(ready);
774 REQUIRE(ready.size() == 1);
775 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
776 auto messageSet = relayer.consumeAllInputsForTimeslice(ready[0].slot);
780 REQUIRE((messageSet[0] |
count_parts{}) == nSplitParts);
784 SECTION(
"SplitPayloadSequence")
787 InputSpec spec1{
"clusters",
"TST",
"COUNTER"};
789 std::vector<InputRoute> inputs = {
793 std::vector<ForwardRoute> forwards;
794 std::vector<InputChannelInfo> infos{1};
796 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
802 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
803 size_t timeslice = 0;
805 std::vector<size_t> sequenceSize;
806 size_t nTotalPayloads = 0;
808 auto createSequence = [&nTotalPayloads, ×lice, &sequenceSize, &transport, &relayer](
size_t nPayloads) ->
void {
809 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
810 std::vector<std::unique_ptr<fair::mq::Message>> messages;
811 messages.reserve(nPayloads + 1);
817 dh.splitPayloadParts = nPayloads;
819 messages.emplace_back(std::move(header));
821 for (
size_t i = 0;
i < nPayloads; ++
i) {
822 messages.emplace_back(transport->CreateMessage(100));
823 *(
reinterpret_cast<size_t*
>(messages.back()->GetData())) = nTotalPayloads;
826 REQUIRE(messages.size() == nPayloads + 1);
828 relayer.relay(messages[0]->GetData(), messages.data(), fakeInfo, messages.size(), nPayloads);
829 sequenceSize.emplace_back(nPayloads);
835 std::vector<RecordAction> ready;
836 relayer.getReadyToProcess(ready);
837 REQUIRE(ready.size() == 1);
838 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
839 auto messageSet = relayer.consumeAllInputsForTimeslice(ready[0].slot);
843 REQUIRE((messageSet[0] |
count_parts{}) == sequenceSize.size());
845 for (
size_t seqid = 0; seqid < sequenceSize.size(); ++seqid) {
849 auto const*
data = (messageSet[0] |
get_payload{seqid, pi})->GetData();
850 REQUIRE(*(
reinterpret_cast<size_t const*
>(
data)) ==
counter);
856 SECTION(
"ProcessDanglingInputs")
858 InputSpec spec{
"condition",
"TST",
"COND"};
859 std::vector<InputRoute> inputs = {
860 InputRoute{spec, 0,
"from_source_to_self", 0}};
862 std::vector<InputChannelInfo> infos{1};
864 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
868 std::vector<fair::mq::Channel>
channels{fair::mq::Channel(
"from_source_to_self")};
869 auto findChannel = [&
channels](std::string
const&
name) -> fair::mq::Channel& {
871 if (ch.GetName() ==
name) {
875 throw std::runtime_error(
"Channel not found: " +
name);
877 proxy.
bind({}, inputs, {}, findChannel, [] {
return false; });
878 ref.registerService(ServiceRegistryHelpers::handleForService<FairMQDeviceProxy>(&proxy));
884 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
885 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
889 dh.splitPayloadIndex = 0;
893 handler.
name =
"test-condition";
895 handler.
lifetime = Lifetime::Condition;
900 for (
size_t si = 0; si <
index.size(); si++) {
902 if (!
index.isValid(slot)) {
904 (
void)
index.setOldestPossibleInput({1}, channelIndex);
917 ref.payload = transport->CreateMessage(4);
920 std::vector<ExpirationHandler> handlers{handler};
921 auto activity = relayer.processDanglingInputs(handlers, {registry},
true);
923 REQUIRE(activity.newSlots == 1);
924 REQUIRE(activity.expiredSlots == 1);
927 std::vector<RecordAction> ready;
928 relayer.getReadyToProcess(ready);
929 REQUIRE(ready.size() == 1);
930 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
932 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
937 SECTION(
"ProcessDanglingInputsSkipsWhenDataPresent")
941 InputSpec spec{
"condition",
"TST",
"COND"};
942 std::vector<InputRoute> inputs = {
943 InputRoute{spec, 0,
"from_source_to_self", 0}};
945 std::vector<InputChannelInfo> infos{1};
947 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
950 std::vector<fair::mq::Channel>
channels{fair::mq::Channel(
"from_source_to_self")};
951 auto findChannel = [&
channels](std::string
const&
name) -> fair::mq::Channel& {
953 if (ch.GetName() ==
name) {
957 throw std::runtime_error(
"Channel not found: " +
name);
959 proxy.
bind({}, inputs, {}, findChannel, [] {
return false; });
960 ref.registerService(ServiceRegistryHelpers::handleForService<FairMQDeviceProxy>(&proxy));
966 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
967 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
971 dh.splitPayloadIndex = 0;
976 handler.
name =
"test-condition";
978 handler.
lifetime = Lifetime::Condition;
981 for (
size_t si = 0; si <
index.size(); si++) {
983 if (!
index.isValid(slot)) {
985 (
void)
index.setOldestPossibleInput({1}, channelIndex);
992 int handlerCallCount = 0;
995 ref.payload = transport->CreateMessage(4);
998 std::vector<ExpirationHandler> handlers{handler};
1001 auto activity1 = relayer.processDanglingInputs(handlers, {registry},
true);
1002 REQUIRE(activity1.expiredSlots == 1);
1003 REQUIRE(handlerCallCount == 1);
1006 auto activity2 = relayer.processDanglingInputs(handlers, {registry},
false);
1007 REQUIRE(activity2.expiredSlots == 0);
1008 REQUIRE(handlerCallCount == 1);
1018 SECTION(
"InterleavedPartsKeepIdentity")
1020 InputSpec spec0{
"clusters",
"TPC",
"CLUSTERS"};
1021 InputSpec spec1{
"its",
"ITS",
"CLUSTERS"};
1022 InputSpec spec2{
"tracks",
"TPC",
"TRACKS"};
1024 std::vector<InputRoute> inputs = {
1030 std::vector<InputChannelInfo> infos{1};
1032 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
1038 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
1039 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
1041 std::array<DataHeader, 3> prototypes;
1042 prototypes[0].dataOrigin =
"TPC";
1043 prototypes[0].dataDescription =
"CLUSTERS";
1044 prototypes[1].dataOrigin =
"ITS";
1045 prototypes[1].dataDescription =
"CLUSTERS";
1046 prototypes[2].dataOrigin =
"TPC";
1047 prototypes[2].dataDescription =
"TRACKS";
1049 auto stampOf = [](
size_t input,
size_t part) -> uint32_t {
1050 return 1000u *
static_cast<uint32_t
>(input + 1) +
static_cast<uint32_t
>(part);
1053 auto relayOne = [&](
size_t input,
size_t part,
size_t timeslice) {
1060 std::array<fair::mq::MessagePtr, 2> msgs;
1062 msgs[1] = transport->CreateMessage(
sizeof(uint32_t));
1063 uint32_t
const stamp = stampOf(input, part);
1064 memcpy(msgs[1]->GetData(), &stamp,
sizeof(stamp));
1066 relayer.relay(msgs[0]->GetData(), msgs.data(), info, 2);
1067 REQUIRE(msgs[0].
get() ==
nullptr);
1068 REQUIRE(msgs[1].
get() ==
nullptr);
1071 std::array<std::pair<size_t, size_t>, 5>
const arrivals = {{{0, 0}, {1, 0}, {0, 1}, {2, 0}, {1, 1}}};
1072 for (
auto const& [input, part] : arrivals) {
1073 relayOne(input, part, 0);
1076 std::vector<RecordAction> ready;
1077 relayer.getReadyToProcess(ready);
1078 REQUIRE(ready.size() == 1);
1079 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
1081 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
1084 std::array<size_t, 3>
const expectedParts = {2, 2, 1};
1085 auto checkContents = [&]() {
1086 for (
size_t i = 0;
i < 3; ++
i) {
1088 for (
size_t p = 0;
p < expectedParts[
i]; ++
p) {
1091 REQUIRE(header.get() !=
nullptr);
1092 REQUIRE(payload.get() !=
nullptr);
1094 memcpy(&seen, payload->GetData(),
sizeof(seen));
1095 REQUIRE(seen == stampOf(
i, p));
1112 SECTION(
"ExpiryDoesNotDisturbNeighbours")
1114 InputSpec dataSpec0{
"clusters",
"TPC",
"CLUSTERS"};
1115 InputSpec condSpec{
"condition",
"TST",
"COND"};
1116 InputSpec dataSpec2{
"tracks",
"TPC",
"TRACKS"};
1118 std::vector<InputRoute> inputs = {
1119 InputRoute{dataSpec0, 0,
"from_source_to_self", 0},
1120 InputRoute{condSpec, 1,
"from_source_to_self", 0},
1121 InputRoute{dataSpec2, 2,
"from_source_to_self", 0},
1124 std::vector<InputChannelInfo> infos{1};
1126 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
1129 std::vector<fair::mq::Channel>
channels{fair::mq::Channel(
"from_source_to_self")};
1130 auto findChannel = [&
channels](std::string
const&
name) -> fair::mq::Channel& {
1132 if (ch.GetName() ==
name) {
1136 throw std::runtime_error(
"Channel not found: " +
name);
1138 proxy.
bind({}, inputs, {}, findChannel, [] {
return false; });
1139 ref.registerService(ServiceRegistryHelpers::handleForService<FairMQDeviceProxy>(&proxy));
1145 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
1146 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
1148 auto stampOf = [](
size_t input) -> uint32_t {
return 7000u +
static_cast<uint32_t
>(input); };
1150 auto relayData = [&](
size_t input,
char const*
origin,
char const*
description) {
1158 std::array<fair::mq::MessagePtr, 2> msgs;
1160 msgs[1] = transport->CreateMessage(
sizeof(uint32_t));
1161 uint32_t
const stamp = stampOf(input);
1162 memcpy(msgs[1]->GetData(), &stamp,
sizeof(stamp));
1164 relayer.relay(msgs[0]->GetData(), msgs.data(), info, 2);
1165 REQUIRE(msgs[0].
get() ==
nullptr);
1170 relayData(0,
"TPC",
"CLUSTERS");
1171 relayData(2,
"TPC",
"TRACKS");
1175 condDh.splitPayloadIndex = 0;
1179 handler.
name =
"test-condition";
1181 handler.
lifetime = Lifetime::Condition;
1190 part.payload = transport->CreateMessage(4);
1193 std::vector<ExpirationHandler> handlers{handler};
1194 auto activity = relayer.processDanglingInputs(handlers, {registry},
true);
1195 REQUIRE(activity.expiredSlots == 1);
1197 std::vector<RecordAction> ready;
1198 relayer.getReadyToProcess(ready);
1199 REQUIRE(ready.size() == 1);
1200 REQUIRE(ready[0].
op == CompletionPolicy::CompletionOp::Consume);
1202 auto result = relayer.consumeAllInputsForTimeslice(ready[0].slot);
1205 for (
size_t i : {0u, 2u}) {
1208 REQUIRE(payload.get() !=
nullptr);
1210 memcpy(&seen, payload->GetData(),
sizeof(seen));
1211 REQUIRE(seen == stampOf(
i));
1219 SECTION(
"RelayAllocationBudget")
1221 constexpr size_t kInputs = 8;
1222 std::vector<InputSpec> specs;
1223 std::vector<InputRoute> inputs;
1224 std::vector<DataHeader> prototypes;
1225 std::array<char const*, kInputs>
const descriptions = {
1226 "CLUSTERS",
"TRACKS",
"DIGITS",
"VERTICES",
"ERRORS",
"CALIB",
"RAWDATA",
"MCLABELS"};
1227 for (
size_t i = 0;
i < kInputs; ++
i) {
1230 specs.emplace_back(
InputSpec{
"in",
"TST", desc});
1232 for (
size_t i = 0;
i < kInputs; ++
i) {
1233 inputs.emplace_back(
InputRoute{specs[
i],
i,
"Fake", 0});
1241 prototypes.push_back(dh);
1244 std::vector<InputChannelInfo> infos{1};
1246 ref.registerService(ServiceRegistryHelpers::handleForService<TimesliceIndex>(&
index));
1252 auto transport = fair::mq::TransportFactory::CreateTransportFactory(
"zeromq");
1253 auto channelAlloc = o2::pmr::getTransportAllocator(transport.get());
1258 auto makeMessages = [&](
size_t timeslice) {
1259 std::vector<std::array<fair::mq::MessagePtr, 2>> msgs(kInputs);
1260 for (
size_t i = 0;
i < kInputs; ++
i) {
1262 msgs[
i][1] = transport->CreateMessage(8);
1267 auto cycle = [&](std::vector<std::array<fair::mq::MessagePtr, 2>>& msgs) {
1268 for (
size_t i = 0;
i < kInputs; ++
i) {
1270 relayer.relay(msgs[
i][0]->GetData(), msgs[
i].data(),
info, 2);
1272 std::vector<RecordAction> ready;
1273 relayer.getReadyToProcess(ready);
1274 REQUIRE(ready.size() == 1);
1275 return relayer.consumeAllInputsForTimeslice(ready[0].slot);
1280 for (
size_t t = 0; t < 4; ++t) {
1281 auto msgs = makeMessages(t);
1282 auto warm = cycle(msgs);
1285 auto msgs = makeMessages(4);
1286 size_t allocations = 0;
1288 AllocationCounter counting;
1289 auto result = cycle(msgs);
1290 allocations = AllocationCounter::count();
1295 REQUIRE(allocations <= 18);