12#ifndef FRAMEWORK_ANALYSIS_TASK_H_
13#define FRAMEWORK_ANALYSIS_TASK_H_
30#include <fairmq/Version.h>
32#include <arrow/compute/kernel.h>
33#include <arrow/table.h>
34#include <gandiva/node.h>
55template <
int64_t BEGIN,
int64_t END,
int64_t STEP = 1>
63static constexpr bool is_enumeration_v =
false;
65template <
int64_t BEGIN,
int64_t END,
int64_t STEP>
66static constexpr bool is_enumeration_v<Enumeration<BEGIN, END, STEP>> =
true;
78struct AnalysisDataProcessorBuilder {
85 if constexpr (soa::relatedByIndex<std::decay_t<G>, std::decay_t<As>>()) {
86 Entry e{soa::getLabelFromTypeForKey<std::decay_t<As>>(
key), soa::getMatcherFromTypeForKey<std::decay_t<As>>(
key),
key,
enabled};
95 }(framework::pack<Args...>{}, bk, bku,
enabled);
98 template <soa::TableRef R>
101 auto spec = soa::tableRef2InputSpec<R>(newOrigin);
102 if (
R.origin_hash !=
"AOD"_h) {
103 spec.metadata.emplace_back(ConfigParamSpec{
"aod-origin-replaced",
VariantType::Bool,
true, {
"\"\""}});
108 auto locate = std::ranges::find_if(iInfos, [&
hash](
auto const& info) {
return info.hash ==
hash; });
109 if (locate == iInfos.end()) {
110 iInfos.emplace_back(
hash, std::vector{std::pair{ai, matcher}});
112 if (std::ranges::none_of(locate->matchers, [&ai, &matcher](
auto const&
match) { return (match.first == ai) && (match.second == matcher); })) {
113 locate->matchers.emplace_back(std::pair{ai, matcher});
119 template <soa::is_table A>
121 static void addExpression(
int, uint32_t, std::vector<ExpressionInfo>&)
125 template <soa::is_filtered_table A>
126 static void addExpression(
int ai, uint32_t
hash, std::vector<ExpressionInfo>& eInfos)
129 eInfos.emplace_back(ai,
hash, std::decay_t<A>::hashes(), std::make_shared<arrow::Schema>(
fields));
132 template <soa::is_iterator A>
133 static void addExpression(
int ai, uint32_t
hash, std::vector<ExpressionInfo>& eInfos)
135 addExpression<typename std::decay_t<A>::parent_t>(ai,
hash, eInfos);
139 template <soa::is_table A>
142 [&
name, &
value, &inputs, &iInfos, &ai, &
hash, newOrigin = std::move(newOrigin)]<
size_t N, std::array<soa::TableRef, N> refs,
size_t... Is>(std::index_sequence<Is...>)
mutable {
143 (addOriginalRef<refs[Is]>(
name,
value, inputs, iInfos, ai,
hash, newOrigin), ...);
144 }.template operator()<A::originals.size(), std::decay_t<A>::originals>(std::make_index_sequence<std::decay_t<A>::originals.size()>());
149 static void addInputsAndExpressions(uint32_t
hash,
const char*
name,
bool value, std::vector<InputSpec>& inputs, std::vector<ExpressionInfo>& eInfos, std::vector<InputInfo>& iInfos,
header::DataOrigin&& newOrigin =
header::DataOrigin{
"AOD"})
152 ([&ai, &
hash, &eInfos, &
name, &
value, &inputs, &iInfos, newOrigin]()
mutable {
154 using T = std::decay_t<As>;
155 addExpression<T>(ai,
hash, eInfos);
156 addInput<T>(
name,
value, inputs, iInfos, ai,
hash, std::move(newOrigin));
162 template <
typename T>
163 inline static bool requestInputsFromArgs(T&, std::string
const&, std::vector<InputSpec>&, std::vector<ExpressionInfo>&, std::vector<InputInfo>&,
header::DataOrigin)
167 template <is_process_configurable T>
168 inline static bool requestInputsFromArgs(T& pc, std::string
const&
name, std::vector<InputSpec>& inputs, std::vector<ExpressionInfo>& eis, std::vector<InputInfo>& iifs,
header::DataOrigin newOrigin =
header::DataOrigin{
"AOD"})
170 AnalysisDataProcessorBuilder::inputsFromArgs(pc.process, (
name +
"/" + pc.name).c_str(), pc.value, inputs, eis, iifs, newOrigin);
173 template <
typename T>
174 inline static bool requestCacheFromArgs(T&,
Cache&,
Cache&)
178 template <is_process_configurable T>
179 inline static bool requestCacheFromArgs(T& pc,
Cache& bk,
Cache& bku)
181 AnalysisDataProcessorBuilder::cacheFromArgs(pc.process, pc.value, bk, bku);
185 template <
typename C, is_enumeration A>
186 static void inputsFromArgs(
void (C::*)(
A),
const char* ,
bool , std::vector<InputSpec>& inputs, std::vector<ExpressionInfo>&, std::vector<InputInfo>&,
header::DataOrigin)
188 std::vector<ConfigParamSpec> inputMetadata;
195 static void inputsFromArgs(
void (C::*)(
A, Args...),
const char*
name,
bool value, std::vector<InputSpec>& inputs, std::vector<ExpressionInfo>& eInfos, std::vector<InputInfo>& iInfos,
header::DataOrigin newOrigin =
header::DataOrigin{
"AOD"})
196 requires(std::is_lvalue_reference_v<A> && (std::is_lvalue_reference_v<Args> && ...))
199 addInputsAndExpressions<typename std::decay_t<A>::parent_t, Args...>(
hash,
name,
value, inputs, eInfos, iInfos, std::move(newOrigin));
204 static void inputsFromArgs(
void (C::*)(Args...),
const char*
name,
bool value, std::vector<InputSpec>& inputs, std::vector<ExpressionInfo>& eInfos, std::vector<InputInfo>& iInfos,
header::DataOrigin newOrigin =
header::DataOrigin{
"AOD"})
205 requires(std::is_lvalue_reference_v<Args> && ...)
208 addInputsAndExpressions<Args...>(
hash,
name,
value, inputs, eInfos, iInfos, std::move(newOrigin));
212 template <
typename C, is_enumeration A>
213 static void cacheFromArgs(
void (C::*)(
A), bool,
Cache&,
Cache&)
218 static void cacheFromArgs(
void (C::*)(
A, Args...),
bool value,
Cache& bk,
Cache& bku)
220 addGroupingCandidates<
A, Args...>(bk, bku,
value);
224 static void cacheFromArgs(
void (C::*)(
A, Args...), bool,
Cache&,
Cache&)
228 template <std::ranges::input_range R>
229 static auto extractTablesFromRecord(InputRecord& record,
R matchers)
231 std::vector<soa::ArrowTableRef> tables;
232 std::ranges::transform(
matchers, std::back_inserter(tables), [&record](
auto const&
m) {
233 return record.get<TableConsumer>(
m.second)->asArrowTable();
238 template <soa::is_table T, std::ranges::input_range R>
239 static auto extractFromRecord(InputRecord& record,
R matchers)
241 return T{extractTablesFromRecord(record,
matchers)};
244 template <soa::is_iterator T, std::ranges::input_range R>
245 static auto extractFromRecord(InputRecord& record,
R matchers)
247 return typename T::parent_t{extractTablesFromRecord(record,
matchers)};
250 template <soa::is_filtered T, std::ranges::input_range R>
256 if (info.selection ==
nullptr) {
261 return typename T::parent_t({table}, info.selection);
263 return T({table}, info.selection);
267 template <is_enumeration T,
int AI, std::ranges::input_range R>
268 static auto extract(InputRecord&,
R, std::vector<ExpressionInfo>&,
size_t)
273 template <soa::is_table_or_iterator T,
int AI, std::ranges::input_range R>
274 static auto extract(InputRecord& record,
R matchers, std::vector<ExpressionInfo>& infos,
size_t phash)
278 return extractFilteredFromRecord<T>(record,
matchers, *std::ranges::find_if(infos, [&phash](
ExpressionInfo const&
i) {
return (
i.processHash == phash &&
i.argumentIndex == AI); }));
280 return extractFromRecord<T>(record,
matchers);
285 static auto bindGroupingTable(InputRecord& record,
R matchers,
void (C::*)(Grouping, Args...), std::vector<ExpressionInfo>& infos)
286 requires(!std::same_as<Grouping, void>)
289 return extract<std::decay_t<Grouping>, 0>(record,
matchers | std::views::filter([](
auto const& pair) {
return pair.first == 0; }), infos,
hash);
293 static auto bindAssociatedTables(InputRecord& record,
R matchers,
void (C::*)(Grouping, Args...), std::vector<ExpressionInfo>& infos)
294 requires(!std::same_as<Grouping, void> &&
sizeof...(Args) > 0)
297 return std::make_tuple(extract<std::decay_t<Args>, has_type_at_v<Args>(pack<Args...>{}) + 1>(record,
matchers | std::views::filter([](
auto const& pair) {
return pair.first == has_type_at_v<Args>(pack<Args...>{}) + 1; }), infos,
hash)...);
301 static void overwriteInternalIndices(std::tuple<As...>& dest, std::tuple<As...>
const&
src)
303 (std::get<As>(dest).bindInternalIndicesTo(&std::get<As>(
src)), ...);
307#if (FAIRMQ_VERSION_DEC >= 111000)
308 static void invokeProcess(Task& task, InputRecord& inputs,
R matchers, PointerReconstructor
const& pointerReconstructor,
void (Task::*processingFunction)(Grouping, Associated...), std::vector<ExpressionInfo>& infos, ArrowTableSlicingCache& slices,
header::DataOrigin newOrigin =
header::DataOrigin{
"AOD"})
310 static void invokeProcess(Task& task, InputRecord& inputs,
R matchers,
void (Task::*processingFunction)(Grouping, Associated...), std::vector<ExpressionInfo>& infos, ArrowTableSlicingCache& slices,
header::DataOrigin newOrigin =
header::DataOrigin{
"AOD"})
313 using G = std::decay_t<Grouping>;
314 auto groupingTable = AnalysisDataProcessorBuilder::bindGroupingTable(inputs,
matchers, processingFunction, infos);
315#if (FAIRMQ_VERSION_DEC >= 111000)
317 groupingTable.setPointerReconstructor(pointerReconstructor);
320 constexpr const int numElements = nested_brace_constructible_size<false, std::decay_t<Task>>() / 10;
323 homogeneous_apply_refs_sized<numElements>([&groupingTable](
auto&
element) {
330 if constexpr (
sizeof...(Associated) == 0) {
332 homogeneous_apply_refs_sized<numElements>([&groupingTable](
auto&
element) {
339 for (
auto&
element : groupingTable) {
340 std::invoke(processingFunction, task, *
element);
343 std::invoke(processingFunction, task, groupingTable);
347 auto associatedTables = AnalysisDataProcessorBuilder::bindAssociatedTables(inputs,
matchers, processingFunction, infos);
350 [&task](
auto&... t)
mutable {
351 (homogeneous_apply_refs_sized<numElements>(
360#if (FAIRMQ_VERSION_DEC >= 111000)
361 std::apply([&pointerReconstructor](
auto&... table) {
362 (table.setPointerReconstructor(pointerReconstructor), ...);
367 auto binder = [&task, &groupingTable, &associatedTables](
auto&
x)
mutable {
368 x.bindExternalIndices(&groupingTable, &std::get<std::decay_t<Associated>>(associatedTables)...);
369 homogeneous_apply_refs_sized<numElements>([&
x](
auto& t)
mutable {
376 groupingTable.bindExternalIndices(&std::get<std::decay_t<Associated>>(associatedTables)...);
380 [&binder](
auto&...
x)
mutable {
386 homogeneous_apply_refs_sized<numElements>([&groupingTable, &associatedTables](
auto& t) {
391 overwriteInternalIndices(associatedTables, associatedTables);
393 auto slicer = GroupSlicer(groupingTable, associatedTables, slices, newOrigin);
394 for (
auto& slice : slicer) {
395 auto associatedSlices = slice.associatedTables();
396#if (FAIRMQ_VERSION_DEC >= 111000)
397 std::apply([&pointerReconstructor](
auto&... table) {
398 (table.setPointerReconstructor(pointerReconstructor), ...);
402 overwriteInternalIndices(associatedSlices, associatedTables);
404 [&binder](
auto&...
x)
mutable {
410 homogeneous_apply_refs_sized<numElements>([&groupingTable](
auto&
x) {
416 [](Task& task,
void (Task::*processingFunction)(Grouping, Associated...), Grouping
g, std::tuple<std::decay_t<Associated>...>& at) {
417 std::invoke(processingFunction, task,
g, std::get<std::decay_t<Associated>>(at)...);
418 }(task, processingFunction, slice.groupingElement(), associatedSlices);
422 homogeneous_apply_refs_sized<numElements>([&groupingTable](
auto&
x) {
428 [](Task& task,
void (Task::*processingFunction)(Grouping, Associated...), Grouping
g, std::tuple<std::decay_t<Associated>...>& at) {
429 std::invoke(processingFunction, task,
g, std::get<std::decay_t<Associated>>(at)...);
430 }(task, processingFunction, groupingTable, associatedTables);
438 std::vector<std::pair<std::string, bool>>
map;
449template <
typename T,
typename...
A>
452 auto task = std::make_shared<T>(std::forward<A>(args)...);
453 for (
auto& setting : second.
map) {
456 return analysis_task_parsers::setProcessSwitch(setting,
element);
460 outputName =
first.value;
464template <
typename T,
typename...
A>
467 auto task = std::make_shared<T>(std::forward<A>(args)...);
468 for (
auto& setting :
first.map) {
471 return analysis_task_parsers::setProcessSwitch(setting,
element);
475 outputName = second.
value;
479template <
typename T,
typename...
A>
480auto getTaskNameSetProcesses(std::string& outputName, SetDefaultProcesses
first,
A... args)
482 auto task = std::make_shared<T>(std::forward<A>(args)...);
483 for (
auto& setting :
first.map) {
486 return analysis_task_parsers::setProcessSwitch(setting,
element);
490 auto type_name_str = type_name<T>();
495template <
typename T,
typename...
A>
496auto getTaskNameSetProcesses(std::string& outputName,
TaskName first,
A... args)
498 auto task = std::make_shared<T>(std::forward<A>(args)...);
499 outputName =
first.value;
503template <
typename T,
typename...
A>
504auto getTaskNameSetProcesses(std::string& outputName,
A... args)
506 auto task = std::make_shared<T>(std::forward<A>(args)...);
507 auto type_name_str = type_name<T>();
515template <
typename T,
typename... Args>
518 TH1::AddDirectory(
false);
520 std::string name_str;
521 auto task = getTaskNameSetProcesses<T>(name_str, args...);
523 auto suffix = ctx.options().get<std::string>(
"workflow-suffix");
524 if (!suffix.empty()) {
527 const char*
name = name_str.c_str();
531 std::vector<OutputSpec> outputs;
532 std::vector<InputSpec> inputs;
533 std::vector<ConfigParamSpec> options;
534 std::vector<ExpressionInfo> expressionInfos;
535 std::vector<InputInfo> inputInfos;
537 std::string newOriginStr;
539 if (ctx.options().hasOption(
"aod-origin-replace")) {
540 newOriginStr = ctx.options().get<std::string>(
"aod-origin-replace");
541 if (newOriginStr.size() > 4UL) {
545 if (!newOriginStr.empty()) {
546 newOrigin.runtimeInit(newOriginStr.c_str(), std::min(newOriginStr.size(), 4UL));
549 constexpr const int numElements = nested_brace_constructible_size<false, std::decay_t<T>>() / 10;
557 if constexpr (
requires { &T::process; }) {
558 AnalysisDataProcessorBuilder::inputsFromArgs(&T::process,
"default",
true, inputs, expressionInfos, inputInfos, newOrigin);
560 homogeneous_apply_refs_sized<numElements>(
561 [
name = name_str, &expressionInfos, &inputs, &inputInfos, &newOrigin](
auto&
x)
mutable {
563 return AnalysisDataProcessorBuilder::requestInputsFromArgs(
x,
name, inputs, expressionInfos, inputInfos, newOrigin);
569 homogeneous_apply_refs_sized<numElements>([&inputs, &newOrigin](
auto&
element) {
575 if (inputs.empty() ==
true) {
576 LOG(warn) <<
"Task " << name_str <<
" has no inputs";
585 for (
auto& input : inputs) {
586 for (
auto& meta : input.metadata) {
587 if (meta.name.starts_with(
"ccdb:") && meta.name !=
"ccdb:") {
599 requiredServices.insert(requiredServices.end(), arrowServices.begin(), arrowServices.end());
607 [task = task, expressionInfos, inputInfos, newOrigin, newOriginStr](
InitContext& ic)
mutable {
609 Cache bindingsKeysUnsorted;
618 homogeneous_apply_refs_sized<numElements>([&eosContext](
auto&
element) {
629 if constexpr (
requires { task->init(ic); }) {
634 homogeneous_apply_refs_sized<numElements>(
638 homogeneous_apply_refs_sized<numElements>([&expressionInfos](
auto&
element) {
644 if constexpr (
requires { &T::process; }) {
645 AnalysisDataProcessorBuilder::cacheFromArgs(&T::process,
true, bindingsKeys, bindingsKeysUnsorted);
647 homogeneous_apply_refs_sized<numElements>(
648 [&bindingsKeys, &bindingsKeysUnsorted](
auto&
x) {
649 return AnalysisDataProcessorBuilder::requestCacheFromArgs(
x, bindingsKeys, bindingsKeysUnsorted);
654 std::ranges::transform(bindingsKeys, bindingsKeys.begin(), [&newOrigin](
Entry&
entry) {
655 if ((entry.matcher.origin == header::DataOrigin{
"AOD"}) && (newOrigin !=
header::DataOrigin{
"AOD"})) {
660 std::ranges::transform(bindingsKeysUnsorted, bindingsKeysUnsorted.begin(), [&newOrigin](
Entry&
entry) {
661 if ((entry.matcher.origin == header::DataOrigin{
"AOD"}) && (newOrigin !=
header::DataOrigin{
"AOD"})) {
670#if (FAIRMQ_VERSION_DEC >= 111000)
671 PointerReconstructor pointerReconstructor(
nullptr);
674 return [task, expressionInfos, inputInfos, newOrigin, hasCCDBTables, pointerReconstructor](
ProcessingContext& pc)
mutable {
675 if (hasCCDBTables && (!pointerReconstructor)) {
678 pointerReconstructor = proxy.getShmPointerReconstructor(spec, 0);
681 return [task, expressionInfos, inputInfos, newOrigin](
ProcessingContext& pc)
mutable {
688 std::ranges::for_each(expressionInfos, [](
auto& info) { info.resetSelection =
true; });
691 homogeneous_apply_refs_sized<numElements>([&slices](
auto&
element) {
700 if constexpr (
requires { task->run(pc); }) {
704 if constexpr (
requires { &T::process; }) {
705 auto loc = std::ranges::find_if(inputInfos, [](
auto const& info) {
return info.hash == o2::framework::TypeIdHelpers::uniqueId<decltype(&T::process)>(); });
706 auto matchers = loc == inputInfos.end() ? std::vector<std::pair<int, ConcreteDataMatcher>>{} : loc->matchers;
707#if (FAIRMQ_VERSION_DEC >= 111000)
708 AnalysisDataProcessorBuilder::invokeProcess(*(task.get()), pc.inputs(),
matchers, pointerReconstructor, &T::process, expressionInfos, slices, newOrigin);
710 AnalysisDataProcessorBuilder::invokeProcess(*(task.get()), pc.inputs(),
matchers, &T::process, expressionInfos, slices, newOrigin);
714 homogeneous_apply_refs_sized<numElements>(
715#
if (FAIRMQ_VERSION_DEC >= 111000)
716 [&pc, &expressionInfos, &task, &slices, &inputInfos, &newOrigin, &pointerReconstructor](
auto&
x) {
718 [&pc, &expressionInfos, &task, &slices, &inputInfos, &newOrigin](
auto&
x) {
721 if (
x.value ==
true) {
723 auto matchers = loc == inputInfos.end() ? std::vector<std::pair<int, ConcreteDataMatcher>>{} : loc->matchers;
724#if (FAIRMQ_VERSION_DEC >= 111000)
725 AnalysisDataProcessorBuilder::invokeProcess(*task.get(), pc.inputs(),
matchers, pointerReconstructor,
x.process, expressionInfos, slices, newOrigin);
727 AnalysisDataProcessorBuilder::invokeProcess(*task.get(), pc.inputs(),
matchers,
x.process, expressionInfos, slices, newOrigin);
std::vector< framework::ConcreteDataMatcher > matchers
std::vector< std::shared_ptr< arrow::Field > > fields
constexpr uint32_t runtime_hash(char const *str)
bool match(const std::vector< std::string > &queries, const char *pattern)
GLuint const GLchar * name
GLenum GLenum GLsizei const GLuint GLboolean enabled
GLsizei const GLfloat * value
typedef void(APIENTRYP PFNGLCULLFACEPROC)(GLenum mode)
bool prepareService(InitContext &, T &)
void bindExternalIndicesPartition(P &, T *...)
void setGroupedCombination(C &, TG &, Ts &...)
Combinations handling.
bool initializeCache(ProcessingContext &, T &)
Cache handling.
bool replaceOrigin(T &, header::DataOrigin const &)
Preslice handling.
bool requestInputs(std::vector< InputSpec > &, T &, header::DataOrigin)
void setPartition(P &, T &...)
bool prepareDelayedOutput(ProcessingContext &, T &)
bool prepareOption(InitContext &, O &)
bool finalizeOutput(ProcessingContext &, T &)
bool registerCache(T &, Cache &, Cache &)
bool postRunOutput(EndOfStreamContext &, T &)
bool createExpressionTrees(std::vector< ExpressionInfo > &, T &)
bool newDataframePartition(T &)
bool prepareOutput(ProcessingContext &, T &)
bool appendOption(std::vector< ConfigParamSpec > &, O &)
Options handling.
bool updateSliceInfo(T &, ArrowTableSlicingCache &)
void bindInternalIndicesPartition(P &, T *)
bool appendCondition(std::vector< InputSpec > &, C &)
Conditions handling.
bool postRunService(EndOfStreamContext &, T &)
bool addService(std::vector< ServiceSpec > &, T &)
Service handling.
bool updatePlaceholders(InitContext &, T &)
Filter handling.
constexpr bool appendOutput(std::vector< OutputSpec > &, T &, uint32_t)
Outputs handling.
bool newDataframeCondition(InputRecord &, C &)
bool updateOutputSpec(T &, header::DataOrigin)
void updateFilterInfo(ExpressionInfo &info, std::shared_ptr< arrow::Table > &table)
Defining ITS Vertex explicitly as messageable.
void updatePairList(Cache &list, Entry &entry)
DataProcessorSpec adaptAnalysisTask(ConfigContext const &ctx, Args &&... args)
std::vector< Entry > Cache
std::string type_to_task_name(std::string_view const &camelCase)
Convert a CamelCase task struct name to snake-case task name.
@ Me
Only quit this data processor.
constexpr auto homogeneous_apply_refs(L l, T &&object)
ConfigParamSpec replaceOrigin(ConfigParamSpec &source, std::string const &originStr)
void wrongOriginReplacement(std::string_view replacement)
std::string cutString(std::string &&str)
auto createFieldsFromColumns(framework::pack< C... >)
void missingFilterDeclaration(int hash, int ai)
std::function< ProcessCallback(InitContext &)> InitCallback
static std::vector< ServiceSpec > defaultServices(std::string extraPlugins="", int numWorkers=0)
Split a string into a vector of strings using : as a separator.
static std::vector< ServiceSpec > arrowServices()
static void addOptionIfMissing(std::vector< ConfigParamSpec > &specs, const ConfigParamSpec &spec)
static ConcreteDataMatcher asConcreteDataMatcher(InputSpec const &input)
static void updateInputList(std::vector< InputSpec > &list, InputSpec &&input)
Updates list of InputSpecs by merging metadata.
static constexpr int64_t step
static constexpr int64_t begin
std::vector< std::pair< std::string, bool > > map
Struct to differentiate task names from possible task string arguments.
TaskName(std::string name)
static constexpr uint32_t uniqueId()
static o2::soa::ArrowTableRef joinTables(std::vector< std::shared_ptr< arrow::Table > > &&tables)
LOG(info)<< "Compressed in "<< sw.CpuTime()<< " s"