11#ifndef O2_FRAMEWORK_DATASPECVIEWS_H_
12#define O2_FRAMEWORK_DATASPECVIEWS_H_
14#include <fairmq/FwdDecls.h>
15#include <fairmq/Message.h>
30 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
35 while (mi <
r.size()) {
36 auto* header = o2::header::get<o2::header::DataHeader*>(
r[mi]->GetData());
38 throw std::runtime_error(
"Not a DataHeader");
40 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
41 count += header->splitPayloadParts;
42 mi += header->splitPayloadParts + 1;
44 count += header->splitPayloadParts ? header->splitPayloadParts : 1;
45 mi += header->splitPayloadParts ? 2 * header->splitPayloadParts : 2;
61 if constexpr (
requires {
r.numInputs(); }) {
72 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
77 while (mi <
r.size()) {
78 auto* header = o2::header::get<o2::header::DataHeader*>(
r[mi]->GetData());
79 auto* sih = o2::header::get<o2::framework::SourceInfoHeader*>(
r[mi]->GetData());
80 auto* dih = o2::header::get<o2::framework::DomainInfoHeader*>(
r[mi]->GetData());
81 if (!header && !sih && !dih) {
82 throw std::runtime_error(
"Header information not found");
89 }
else if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
91 mi += header->splitPayloadParts + 1;
93 count += header->splitPayloadParts ? header->splitPayloadParts : 1;
94 mi += header->splitPayloadParts ? 2 * header->splitPayloadParts : 2;
105 template <
typename R>
106 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
111 while (mi <
r.size()) {
112 auto* header = o2::header::get<o2::header::DataHeader*>(
r[mi]->GetData());
114 throw std::runtime_error(
"Not a DataHeader");
117 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
119 count += header->splitPayloadParts;
121 return {mi, mi + 1 + diff};
123 mi += header->splitPayloadParts + 1;
124 }
else if (header->splitPayloadParts > 1 && header->splitPayloadIndex != header->splitPayloadParts) {
127 if (diff < header->splitPayloadParts) {
128 return {mi + 2 * diff, mi + 2 * diff + 1};
130 count += header->splitPayloadParts;
131 mi += 2 * header->splitPayloadParts;
141 throw std::runtime_error(
"Payload not found");
158 template <
typename R>
159 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
163 auto* header = o2::header::get<o2::header::DataHeader*>(
r[hIdx]->GetData());
165 throw std::runtime_error(
"Not a DataHeader");
167 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
174 size_t nextHIdx = hIdx + header->splitPayloadParts + 1;
175 return {nextHIdx, nextHIdx + 1};
178 return {hIdx + 2, hIdx + 3};
186 template <
typename R>
187 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
192 while (mi <
r.size()) {
193 auto* header = o2::header::get<o2::header::DataHeader*>(
r[mi]->GetData());
195 throw std::runtime_error(
"Not a DataHeader");
197 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
199 return {mi, mi + 1 + self.
subPart};
202 mi += header->splitPayloadParts + 1;
205 return {mi, mi + self.
subPart + 1};
211 throw std::runtime_error(
"Payload not found");
218 template <
typename R>
219 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
230 template <
typename R>
231 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
242 template <
typename R>
243 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
249 while (mi <
r.size()) {
250 auto* header = o2::header::get<o2::header::DataHeader*>(
r[mi]->GetData());
252 throw std::runtime_error(
"Not a DataHeader");
254 if (header->splitPayloadParts > 1 && (header->splitPayloadIndex == header->splitPayloadParts)) {
258 return header->splitPayloadParts;
262 mi += header->splitPayloadParts + 1;
268 auto pairs = header->splitPayloadParts ? header->splitPayloadParts : 1;
269 if (self.
n <
count + pairs) {
282 template <
typename R>
283 requires requires(
R r) {
requires std::ranges::random_access_range<
decltype(
r.sets)>; }
286 return std::span(
r.sets[self.
slot.
index *
r.inputsPerSlot]);
292 template <
typename R>
293 requires std::ranges::random_access_range<R>
Defining ITS Vertex explicitly as messageable.
friend size_t operator|(R &&r, count_parts self)
friend size_t operator|(R &&r, count_payloads self)
friend DataRefIndices operator|(R &&r, get_dataref_indices self)
friend DataRefIndices operator|(R &&r, get_next_pair self)
friend size_t operator|(R &&r, get_num_payloads self)
friend DataRefIndices operator|(R &&r, get_pair self)
friend auto & operator|(R &&r, get_payload self)