Skip to content
1 change: 1 addition & 0 deletions Framework/AnalysisSupport/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ endif()
o2_add_library(FrameworkOnDemandTablesSupport
SOURCES src/OnDemandPlugin.cxx
src/AODReaderHelpers.cxx
src/AODSliceHelpers.cxx
PRIVATE_INCLUDE_DIRECTORIES ${CMAKE_CURRENT_LIST_DIR}/src
PUBLIC_LINK_LIBRARIES O2::Framework ${EXTRA_TARGETS})

Expand Down
66 changes: 66 additions & 0 deletions Framework/AnalysisSupport/src/AODSliceHelpers.cxx
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
// Copyright 2019-2020 CERN and copyright holders of ALICE O2.
// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
// All rights not expressly granted are reserved.
//
// This software is distributed under the terms of the GNU General Public
// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
//
// In applying this license CERN does not waive the privileges and immunities
// granted to it by virtue of its status as an Intergovernmental Organization
// or submit itself to any jurisdiction.

#include "AODSliceHelpers.h"

#include "Framework/ArrowTableSlicingCache.h"
#include "Framework/ConfigParamRegistry.h"
#include "Framework/DanglingEdgesContext.h"
#include "Framework/DataAllocator.h"
#include "Framework/DataSpecUtils.h"
#include "Framework/InputRecord.h"
#include "Framework/TableConsumer.h"

namespace o2::framework::helpers
{
namespace
{
Entry sourceEntry(InputSpec const& spec)
{
auto source = DataSpecUtils::fromMetadataString(std::ranges::find_if(spec.metadata, [](ConfigParamSpec const& cps) { return cps.name.starts_with("slice-source"); })->defaultValue.get<std::string>());
return {source.binding, DataSpecUtils::asConcreteDataMatcher(source), std::ranges::find_if(spec.metadata, [](ConfigParamSpec const& cps) { return cps.name.starts_with("slice-key"); })->defaultValue.get<std::string>()};
}

struct Sliceable {
Entry entry;
ConcreteDataMatcher output;
bool sorted;

explicit Sliceable(InputSpec const& spec)
: entry{sourceEntry(spec)},
output{DataSpecUtils::asConcreteDataMatcher(spec)},
sorted{std::ranges::find_if(spec.metadata, [](ConfigParamSpec const& cps) { return cps.name.starts_with("sorted"); })->defaultValue.get<bool>()}
{
}

std::shared_ptr<arrow::Table> materialize(ProcessingContext& pc) const
{
auto source = pc.inputs().get<TableConsumer>(entry.matcher)->asArrowTable();
return sorted ? SliceInfo::makeSorted(entry, source) : SliceInfo::makeUnsorted(entry, source);
}
};
} // namespace

AlgorithmSpec AODSliceHelpers::arrowTablesSlicerCallback(ConfigContext const& /*ctx*/)
{
return AlgorithmSpec::InitCallback{[](InitContext& ic) {
// each slicer handles the group of slice infos for the tables from a single provider
auto const& requested = ic.services().get<DanglingEdgesContext>().slicerGroups[ic.options().get<int>("slicer-group")];
std::vector<Sliceable> sliceables;
sliceables.reserve(requested.size());
std::ranges::transform(requested, std::back_inserter(sliceables), [](auto const& i) { return Sliceable{i}; });
return [sliceables](ProcessingContext& pc) {
auto outputs = pc.outputs();
std::ranges::for_each(sliceables, [&pc, &outputs](auto const& sliceable) { outputs.adopt(Output{sliceable.output.origin, sliceable.output.description, sliceable.output.subSpec}, sliceable.materialize(pc)); });
};
}};
}
} // namespace o2::framework::helpers
25 changes: 25 additions & 0 deletions Framework/AnalysisSupport/src/AODSliceHelpers.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
// Copyright 2019-2020 CERN and copyright holders of ALICE O2.
// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
// All rights not expressly granted are reserved.
//
// This software is distributed under the terms of the GNU General Public
// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
//
// In applying this license CERN does not waive the privileges and immunities
// granted to it by virtue of its status as an Intergovernmental Organization
// or submit itself to any jurisdiction.

#ifndef AODSLICEHELPERS_H
#define AODSLICEHELPERS_H

#include "Framework/AlgorithmSpec.h"
namespace o2::framework::helpers
{

struct AODSliceHelpers {
static AlgorithmSpec arrowTablesSlicerCallback(ConfigContext const& /*ctx*/);
};

} // namespace o2::framework::helpers

#endif // AODSLICEHELPERS_H
9 changes: 9 additions & 0 deletions Framework/AnalysisSupport/src/OnDemandPlugin.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
#include "Framework/Plugins.h"
#include "Framework/AlgorithmSpec.h"
#include "AODReaderHelpers.h"
#include "AODSliceHelpers.h"

struct ExtendedTableSpawner : o2::framework::AlgorithmPlugin {
o2::framework::AlgorithmSpec create(o2::framework::ConfigContext const& config) override
Expand All @@ -26,7 +27,15 @@ struct IndexTableBuilder : o2::framework::AlgorithmPlugin {
}
};

struct ArrowTableSlicer : o2::framework::AlgorithmPlugin {
o2::framework::AlgorithmSpec create(o2::framework::ConfigContext const& config) override
{
return o2::framework::helpers::AODSliceHelpers::arrowTablesSlicerCallback(config);
}
};

DEFINE_DPL_PLUGINS_BEGIN
DEFINE_DPL_PLUGIN_INSTANCE(ExtendedTableSpawner, CustomAlgorithm);
DEFINE_DPL_PLUGIN_INSTANCE(IndexTableBuilder, CustomAlgorithm);
DEFINE_DPL_PLUGIN_INSTANCE(ArrowTableSlicer, CustomAlgorithm);
DEFINE_DPL_PLUGINS_END
41 changes: 16 additions & 25 deletions Framework/Core/include/Framework/ASoA.h
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,8 @@
#if defined(__CLING__)
#error "Please do not include this file in ROOT dictionary generation"
#endif
#include "Framework/InputSpec.h"
#include "Framework/Concepts.h"
#include "Framework/ConcreteDataMatcher.h"
#include "Framework/Pack.h" // IWYU pragma: export
#include "Framework/FunctionalHelpers.h" // IWYU pragma: export
#include "Headers/DataHeader.h" // IWYU pragma: export
Expand Down Expand Up @@ -1351,7 +1351,7 @@ static constexpr std::pair<bool, framework::ConcreteDataMatcher> hasKeyM(std::st
}

void notFoundColumn(const char* label, const char* key);
void missingOptionalPreslice(const char* label, const char* key);
void missingPreslice(const char* label, const char* key);

template <with_originals T, bool OPT = false>
static constexpr std::string getLabelFromTypeForKey(std::string_view key)
Expand Down Expand Up @@ -1421,7 +1421,6 @@ namespace o2::framework
/// tracks origin in bindingKey matcher to handle the correct arguments
struct PreslicePolicyBase {
static constexpr void isPreslicePolicy() {};
const std::string binding;
Entry bindingKey;

bool isMissing() const;
Expand All @@ -1448,29 +1447,27 @@ struct PresliceBase : public Policy {
constexpr static bool optional = OPT;
using target_t = T;
using policy_t = Policy;
const std::string binding;

PresliceBase(expressions::BindingNode index_)
: Policy{PreslicePolicyBase{{o2::soa::getLabelFromTypeForKey<T, OPT>(std::string{index_.name})}, Entry(o2::soa::getLabelFromTypeForKey<T, OPT>(std::string{index_.name}), o2::soa::getMatcherFromTypeForKey<T, OPT>(std::string{index_.name}), std::string{index_.name})}, {}}
: Policy{Entry(
o2::soa::getLabelFromTypeForKey<T, true>(std::string{index_.name}),
o2::soa::getMatcherFromTypeForKey<T, true>(std::string{index_.name}),
std::string{index_.name})}
{
}

o2::soa::ArrowTableRef getSliceFor(int value, o2::soa::ArrowTableRef const& input) const
{
if constexpr (OPT) {
if (Policy::isMissing()) {
return {nullptr, {0, 0}};
}
if (Policy::isMissing()) {
return {nullptr, {0, 0}};
}
return Policy::getSliceFor(value, input);
}

std::span<const int64_t> getSliceFor(int value) const
{
if constexpr (OPT) {
if (Policy::isMissing()) {
return {};
}
if (Policy::isMissing()) {
return {};
}
return Policy::getSliceFor(value);
}
Expand Down Expand Up @@ -1526,10 +1523,8 @@ template <typename T, typename C, typename Policy, bool OPT>
requires std::same_as<Policy, framework::PreslicePolicySorted> && (o2::soa::is_binding_compatible_v<C, T>())
auto doSliceBy(T const* table, o2::framework::PresliceBase<C, Policy, OPT> const& container, int value)
{
if constexpr (OPT) {
if (container.isMissing()) {
missingOptionalPreslice(getLabelFromType<std::decay_t<T>>().data(), container.bindingKey.key.c_str());
}
if (container.isMissing()) {
missingPreslice(getLabelFromType<std::decay_t<T>>().data(), container.bindingKey.key.c_str());
}
auto out = container.getSliceFor(value, table->asArrowTableRef());
auto t = typename T::self_t({out});
Expand Down Expand Up @@ -1568,10 +1563,8 @@ template <typename T, typename C, typename Policy, bool OPT>
requires std::same_as<Policy, framework::PreslicePolicyGeneral> && (o2::soa::is_binding_compatible_v<C, T>())
auto doSliceBy(T const* table, o2::framework::PresliceBase<C, Policy, OPT> const& container, int value)
{
if constexpr (OPT) {
if (container.isMissing()) {
missingOptionalPreslice(getLabelFromType<std::decay_t<T>>().data(), container.bindingKey.key.c_str());
}
if (container.isMissing()) {
missingPreslice(getLabelFromType<std::decay_t<T>>().data(), container.bindingKey.key.c_str());
}
auto selection = container.getSliceFor(value);
return doSliceByHelper(table, selection);
Expand Down Expand Up @@ -1601,10 +1594,8 @@ template <soa::is_filtered_table T, typename C, bool OPT>
requires(o2::soa::is_binding_compatible_v<C, T>())
auto doFilteredSliceBy(T const* table, o2::framework::PresliceBase<C, framework::PreslicePolicySorted, OPT> const& container, int value)
{
if constexpr (OPT) {
if (container.isMissing()) {
missingOptionalPreslice(getLabelFromType<T>().data(), container.bindingKey.key.c_str());
}
if (container.isMissing()) {
missingPreslice(getLabelFromType<T>().data(), container.bindingKey.key.c_str());
}
auto slice = container.getSliceFor(value, table->asArrowTableRef());
return prepareFilteredSlice(table, slice);
Expand Down
2 changes: 2 additions & 0 deletions Framework/Core/include/Framework/AnalysisHelpers.h
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@ struct InputInfo {
uint32_t hash;
std::vector<std::pair<int, ConcreteDataMatcher>> matchers;
};

void updateInputInfos(std::vector<InputInfo>& iInfos, ConcreteDataMatcher&& matcher, uint32_t hash, int ai);
} // namespace o2::framework

namespace o2::soa
Expand Down
91 changes: 75 additions & 16 deletions Framework/Core/include/Framework/AnalysisManagers.h
Original file line number Diff line number Diff line change
Expand Up @@ -627,6 +627,73 @@ bool replaceOrigin(T& presliceGroup, header::DataOrigin const& newOrigin)
return true;
}

template <typename T>
requires(!is_preslice<T> && !is_preslice_group<T>)
bool addSlicingInputs(T&, std::vector<InputSpec>&, header::DataOrigin const&)
{
return false;
}

/// check if any of the tables the sliced type is based on is an input of the task
template <soa::is_table T>
bool isSlicedTableInput(std::vector<InputSpec> const& inputs, header::DataOrigin const& newOrigin)
{
auto isInput = [&inputs, &newOrigin](ConcreteDataMatcher matcher) {
if ((matcher.origin == header::DataOrigin{"AOD"}) && (newOrigin != header::DataOrigin{"AOD"})) {
matcher = replaceOrigin(matcher, newOrigin);
}
return std::ranges::any_of(inputs, [&matcher](InputSpec const& input) { return DataSpecUtils::match(input, matcher); });
};
return [&isInput]<size_t... Is>(std::index_sequence<Is...>) {
return (isInput(o2::aod::matcher<T::originals[Is]>()) || ...);
}(std::make_index_sequence<T::originals.size()>{});
}

/// all the process function inputs are already added at this point, so a Preslice can only
/// amend them. Depending on whether the sliced table is an input of the task and whether it has
/// the index column, there are 4 cases:
/// 1. no table, no column - likely an incorrect declaration, warning for both Preslice and PresliceOptional
/// 2. no table, column - Preslice that never works, or a common declaration in a templated task that is
/// not effective in this specialization, warning for Preslice only
/// 3. table, no column - the intended case for PresliceOptional, a mistake for Preslice, warning for Preslice only
/// 4. table, column - slicing input is added
template <is_preslice T>
bool addSlicingInputs(T& preslice, std::vector<InputSpec>& inputs, header::DataOrigin const& newOrigin)
{
using target_t = typename T::target_t;
auto const& [binding, matcher, key, enabled] = preslice.bindingKey;
if (preslice.isMissing()) {
if (!isSlicedTableInput<target_t>(inputs, newOrigin)) {
LOGP(warn, "Preslice declared on {} is skipped: {} is not an input of any process function and does not have column {}, the declaration is likely incorrect",
o2::soa::getLabelFromType<target_t>(), o2::soa::getLabelFromType<target_t>(), key);
} else if constexpr (!T::optional) {
LOGP(warn, "Preslice declared on {} is skipped: it does not have column {}, use PresliceOptional if the column is not always expected",
o2::soa::getLabelFromType<target_t>(), key);
}
return true;
}
if (std::ranges::none_of(inputs, [&matcher](InputSpec const& input) { return DataSpecUtils::match(input, matcher); })) {
if constexpr (!T::optional) {
LOGP(warn, "Preslice declared on {}/{} ({}) is skipped: {} is not an input of any process function, use PresliceOptional if the declaration is not effective in every specialization of a templated task",
binding, key, DataSpecUtils::describe(matcher), binding);
}
return true;
}
DataSpecUtils::updateInputList(inputs, inputForEntry(preslice.bindingKey, std::same_as<typename T::policy_t, framework::PreslicePolicySorted>));
return true;
}

template <is_preslice_group T>
bool addSlicingInputs(T&& presliceGroup, std::vector<InputSpec>& inputs, header::DataOrigin const& newOrigin)
{
homogeneous_apply_refs<true>(
[&inputs, &newOrigin](auto& preslice) {
return addSlicingInputs(preslice, inputs, newOrigin);
},
presliceGroup);
return true;
}

template <typename T>
requires(!is_preslice<T> && !is_preslice_group<T>)
bool registerCache(T&, Cache&, Cache&)
Expand All @@ -638,10 +705,8 @@ template <is_preslice T>
requires std::same_as<typename T::policy_t, framework::PreslicePolicySorted>
bool registerCache(T& preslice, Cache& bsks, Cache&)
{
if constexpr (T::optional) {
if (preslice.binding == "[MISSING]") {
return true;
}
if (preslice.isMissing()) {
return true;
}
auto locate = std::find(bsks.begin(), bsks.end(), preslice.getBindingKey());
if (locate == bsks.end()) {
Expand All @@ -656,10 +721,8 @@ template <is_preslice T>
requires std::same_as<typename T::policy_t, framework::PreslicePolicyGeneral>
bool registerCache(T& preslice, Cache&, Cache& bsksU)
{
if constexpr (T::optional) {
if (preslice.binding == "[MISSING]") {
return true;
}
if (preslice.isMissing()) {
return true;
}
auto locate = std::find(bsksU.begin(), bsksU.end(), preslice.getBindingKey());
if (locate == bsksU.end()) {
Expand Down Expand Up @@ -688,10 +751,8 @@ template <is_preslice T>
static bool updateSliceInfo(T& preslice, ArrowTableSlicingCache& cache)
requires std::same_as<typename T::policy_t, framework::PreslicePolicySorted>
{
if constexpr (T::optional) {
if (preslice.binding == "[MISSING]") {
return true;
}
if (preslice.isMissing()) {
return true;
}
preslice.updateSliceInfo(cache.getCacheFor(preslice.getBindingKey()));
return true;
Expand All @@ -701,10 +762,8 @@ template <is_preslice T>
static bool updateSliceInfo(T& preslice, ArrowTableSlicingCache& cache)
requires std::same_as<typename T::policy_t, framework::PreslicePolicyGeneral>
{
if constexpr (T::optional) {
if (preslice.binding == "[MISSING]") {
return true;
}
if (preslice.isMissing()) {
return true;
}
preslice.updateSliceInfo(cache.getCacheUnsortedFor(preslice.getBindingKey()));
return true;
Expand Down
9 changes: 9 additions & 0 deletions Framework/Core/include/Framework/AnalysisSupportHelpers.h
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,15 @@ struct AnalysisSupportHelpers {
std::vector<InputSpec>& requestedAODs,
std::vector<InputSpec>& requestedDYNs,
DataProcessorSpec& publisher);
static void addMissingOutputsToSlicer(std::vector<InputSpec> const& requestedSLCs,
DataProcessorSpec& publisher);
/// Split the requested slice infos into groups by the device providing the sliced table
/// (the AOD reader, if none of the providers has it) and create a slicer device for each
/// group, so that each slicer depends on a single device and does not create loops.
/// Each slicer is returned together with the name of its provider.
static std::vector<std::pair<std::string, DataProcessorSpec>> makeSlicers(std::vector<InputSpec> const& requestedSLCs,
std::vector<DataProcessorSpec const*> const& providers,
std::vector<std::vector<InputSpec>>& slicerGroups);

/// Match all inputs of kind ATSK and write them to a ROOT file,
/// one root file per originating task.
Expand Down
Loading
Loading