Project
Loading...
Searching...
No Matches
DataModelViews.h
Go to the documentation of this file.
1// Copyright 2019-2025 CERN and copyright holders of ALICE O2.
2// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
3// All rights not expressly granted are reserved.
4//
5// This software is distributed under the terms of the GNU General Public
6// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
7//
8// In applying this license CERN does not waive the privileges and immunities
9// granted to it by virtue of its status as an Intergovernmental Organization
10// or submit itself to any jurisdiction.
11#ifndef O2_FRAMEWORK_DATASPECVIEWS_H_
12#define O2_FRAMEWORK_DATASPECVIEWS_H_
13
14#include <fairmq/FwdDecls.h>
15#include <fairmq/Message.h>
16#include "DomainInfoHeader.h"
17#include "SourceInfoHeader.h"
18#include "Headers/DataHeader.h"
19#include "Framework/DataRef.h"
21#include <ranges>
22#include <span>
23
24namespace o2::framework
25{
26
28 // ends the pipeline, returns the container
29 template <typename R>
30 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
31 friend size_t operator|(R&& r, count_payloads self)
32 {
33 size_t count = 0;
34 size_t mi = 0;
35 while (mi < r.size()) {
36 auto* header = o2::header::get<o2::header::DataHeader*>(r[mi]->GetData());
37 if (!header) {
38 throw std::runtime_error("Not a DataHeader");
39 }
40 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
41 count += header->splitPayloadParts;
42 mi += header->splitPayloadParts + 1;
43 } else {
44 count += header->splitPayloadParts ? header->splitPayloadParts : 1;
45 mi += header->splitPayloadParts ? 2 * header->splitPayloadParts : 2;
46 }
47 }
48 return count;
49 }
50};
51
52// How many inputs a consumed record holds. A record is either a vector of
53// per-input message sets or an arena keeping them in one buffer; both answer
54// this, but they spell it differently, so ask through here and callers stay put
55// when the storage underneath them changes.
57 // ends the pipeline, returns the number of inputs
58 template <typename R>
59 friend size_t operator|(R&& r, count_inputs self)
60 {
61 if constexpr (requires { r.numInputs(); }) {
62 return r.numInputs();
63 } else {
64 return r.size();
65 }
66 }
67};
68
70 // ends the pipeline, returns the number of parts
71 template <typename R>
72 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
73 friend size_t operator|(R&& r, count_parts self)
74 {
75 size_t count = 0;
76 size_t mi = 0;
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");
83 }
84 // We skip oldest possible timeframe / end of stream and not consider it
85 // as actual parts.
86 if (dih || sih) {
87 count += 1;
88 mi += 2;
89 } else if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
90 count += 1;
91 mi += header->splitPayloadParts + 1;
92 } else {
93 count += header->splitPayloadParts ? header->splitPayloadParts : 1;
94 mi += header->splitPayloadParts ? 2 * header->splitPayloadParts : 2;
95 }
96 }
97 return count;
98 }
99};
100
101// DataRefIndices is defined in Framework/DataRef.h
102
103struct get_pair {
104 size_t pairId;
105 template <typename R>
106 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
108 {
109 size_t count = 0;
110 size_t mi = 0;
111 while (mi < r.size()) {
112 auto* header = o2::header::get<o2::header::DataHeader*>(r[mi]->GetData());
113 if (!header) {
114 throw std::runtime_error("Not a DataHeader");
115 }
116 size_t diff = self.pairId - count;
117 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
118 // New style: one header followed by splitPayloadParts contiguous payloads.
119 count += header->splitPayloadParts;
120 if (self.pairId < count) {
121 return {mi, mi + 1 + diff};
122 }
123 mi += header->splitPayloadParts + 1;
124 } else if (header->splitPayloadParts > 1 && header->splitPayloadIndex != header->splitPayloadParts) {
125 // Old style multi-part: splitPayloadParts [header, payload] pairs.
126 // We are at the first pair of the block; jump directly.
127 if (diff < header->splitPayloadParts) {
128 return {mi + 2 * diff, mi + 2 * diff + 1};
129 }
130 count += header->splitPayloadParts;
131 mi += 2 * header->splitPayloadParts;
132 } else {
133 // Single [header, payload] pair (splitPayloadParts == 0).
134 if (self.pairId == count) {
135 return {mi, mi + 1};
136 }
137 count += 1;
138 mi += 2;
139 }
140 }
141 throw std::runtime_error("Payload not found");
142 }
143};
144
145// Advance from a DataRefIndices to the next one in O(1), reading only the
146// current header. Intended for use in iterators so that ++ is O(1) rather
147// than the O(n) while-loop that get_pair requires.
148//
149// New-style block (splitPayloadIndex == splitPayloadParts > 1):
150// layout: [header, payload_0, payload_1, ..., payload_{N-1}]
151// advance within block while payloads remain, then jump to the next block.
152//
153// Old-style block (splitPayloadIndex != splitPayloadParts, splitPayloadParts > 1)
154// or single pair (splitPayloadParts == 0):
155// layout: [header, payload] – always advance by two messages.
158 template <typename R>
159 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
161 {
162 size_t hIdx = self.current.headerIdx;
163 auto* header = o2::header::get<o2::header::DataHeader*>(r[hIdx]->GetData());
164 if (!header) {
165 throw std::runtime_error("Not a DataHeader");
166 }
167 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
168 // New-style block: one header followed by splitPayloadParts contiguous payloads.
169 if (self.current.payloadIdx < hIdx + header->splitPayloadParts) {
170 // More sub-payloads remain in this block.
171 return {hIdx, self.current.payloadIdx + 1};
172 }
173 // Last sub-payload consumed; move to the first pair of the next block.
174 size_t nextHIdx = hIdx + header->splitPayloadParts + 1;
175 return {nextHIdx, nextHIdx + 1};
176 }
177 // Old-style [header, payload] pairs or a single pair: advance by two messages.
178 return {hIdx + 2, hIdx + 3};
179 }
180};
181
183 size_t part;
184 size_t subPart;
185 // ends the pipeline, returns the number of parts
186 template <typename R>
187 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
189 {
190 size_t count = 0;
191 size_t mi = 0;
192 while (mi < r.size()) {
193 auto* header = o2::header::get<o2::header::DataHeader*>(r[mi]->GetData());
194 if (!header) {
195 throw std::runtime_error("Not a DataHeader");
196 }
197 if (header->splitPayloadParts > 1 && header->splitPayloadIndex == header->splitPayloadParts) {
198 if (self.part == count) {
199 return {mi, mi + 1 + self.subPart};
200 }
201 count += 1;
202 mi += header->splitPayloadParts + 1;
203 } else {
204 if (self.part == count) {
205 return {mi, mi + self.subPart + 1};
206 }
207 count += 1;
208 mi += 2;
209 }
210 }
211 throw std::runtime_error("Payload not found");
212 }
213};
214
216 size_t id;
217 // ends the pipeline, returns the number of parts
218 template <typename R>
219 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
220 friend auto& operator|(R&& r, get_header self)
221 {
222 return r[(r | get_dataref_indices{self.id, 0}).headerIdx];
223 }
224};
225
227 size_t part;
228 size_t subPart;
229 // ends the pipeline, returns the number of parts
230 template <typename R>
231 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
232 friend auto& operator|(R&& r, get_payload self)
233 {
234 return r[(r | get_dataref_indices{self.part, self.subPart}).payloadIdx];
235 }
236};
237
239 size_t n;
240 // ends the pipeline, returns the number of payloads which are associated
241 // to the multipart n-th sequence of messages found in the range
242 template <typename R>
243 requires std::ranges::random_access_range<R> && std::ranges::sized_range<R>
244 friend size_t operator|(R&& r, get_num_payloads self)
245 {
246 size_t count = 0;
247 size_t mi = 0;
248 // Un
249 while (mi < r.size()) {
250 auto* header = o2::header::get<o2::header::DataHeader*>(r[mi]->GetData());
251 if (!header) {
252 throw std::runtime_error("Not a DataHeader");
253 }
254 if (header->splitPayloadParts > 1 && (header->splitPayloadIndex == header->splitPayloadParts)) {
255 // This is the case for the new multi payload messages where the number of parts
256 // is as many as the splitPayloadParts number.
257 if (self.n == count) {
258 return header->splitPayloadParts;
259 }
260 // For multipayload we skip all the parts and their associated header
261 count += 1;
262 mi += header->splitPayloadParts + 1;
263 } else {
264 // This is the case of a multipart (header, payload), (header, payload), ...
265 // sequence where we know how many pairs are there.
266 // When splitPayloadParts == 0, it means it is a non-multipart (header, payload)
267 // pair. Each pair has exactly 1 payload.
268 auto pairs = header->splitPayloadParts ? header->splitPayloadParts : 1;
269 if (self.n < count + pairs) {
270 return 1;
271 }
272 count += pairs;
273 mi += 2 * pairs;
274 }
275 }
276 return 0;
277 }
278};
279
282 template <typename R>
283 requires requires(R r) { requires std::ranges::random_access_range<decltype(r.sets)>; }
284 friend auto operator|(R&& r, inputs_for_slot self)
285 {
286 return std::span(r.sets[self.slot.index * r.inputsPerSlot]);
287 }
288};
289
291 size_t inputIdx;
292 template <typename R>
293 requires std::ranges::random_access_range<R>
294 friend std::span<fair::mq::MessagePtr> operator|(R&& r, messages_for_input self)
295 {
296 return std::span(r[self.inputIdx]);
297 }
298};
299
300// FIXME: we should use special index classes in place of size_t
301// FIXME: we need something to substitute a range in the store with another
302
303} // namespace o2::framework
304
305#endif // O2_FRAMEWORK_DATASPECVIEWS_H_
GLint GLsizei count
Definition glcorearb.h:399
GLboolean r
Definition glcorearb.h:1233
Defining ITS Vertex explicitly as messageable.
Definition Cartesian.h:288
friend size_t operator|(R &&r, count_inputs self)
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 auto & operator|(R &&r, get_header 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)
friend auto operator|(R &&r, inputs_for_slot self)
friend std::span< fair::mq::MessagePtr > operator|(R &&r, messages_for_input self)