Skip to content

Commit 96c9952

Browse files
committed
GPU Workflow: Clean up sending of empty dummy messages for missing outputs
1 parent ac90871 commit 96c9952

1 file changed

Lines changed: 11 additions & 8 deletions

File tree

GPU/Workflow/src/GPUWorkflowSpec.cxx

Lines changed: 11 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -97,6 +97,7 @@
9797
#include <sys/stat.h>
9898
#include <fcntl.h>
9999
#include <chrono>
100+
#include <unordered_set>
100101

101102
using namespace o2::framework;
102103
using namespace o2::header;
@@ -678,21 +679,23 @@ void GPURecoWorkflowSpec::run(ProcessingContext& pc)
678679
using outputBufferUninitializedVector = std::decay_t<decltype(pc.outputs().make<DataAllocator::UninitializedVector<outputDataType>>(Output{"", "", 0}))>;
679680
using outputBufferType = std::pair<std::optional<std::reference_wrapper<outputBufferUninitializedVector>>, outputDataType*>;
680681
std::vector<outputBufferType> outputBuffers(GPUInterfaceOutputs::count(), {std::nullopt, nullptr});
682+
std::unordered_set<std::string> outputsCreated;
681683

682-
auto setOutputAllocator = [this, &outputBuffers, &outputRegions, &pc](const char* name, bool condition, GPUOutputControl& region, auto&& outputSpec, size_t offset = 0) {
684+
auto setOutputAllocator = [this, &outputBuffers, &outputRegions, &pc, &outputsCreated](const char* name, bool condition, GPUOutputControl& region, auto&& outputSpec, size_t offset = 0) {
683685
if (condition) {
684686
auto& buffer = outputBuffers[outputRegions.getIndex(region)];
685687
if (mConfParam->allocateOutputOnTheFly) {
686-
region.allocator = [this, name, &buffer, &pc, outputSpec = std::move(outputSpec), offset](size_t size) -> void* {
688+
region.allocator = [this, name, &buffer, &pc, outputSpec = std::move(outputSpec), offset, &outputsCreated](size_t size) -> void* {
687689
size += offset;
688690
if (mVerbosity) {
689-
LOG(info) << "ALLOCATING " << size << " bytes for " << std::get<DataOrigin>(outputSpec).template as<std::string>() << "/" << std::get<DataDescription>(outputSpec).template as<std::string>() << "/" << std::get<2>(outputSpec);
691+
LOG(info) << "ALLOCATING " << size << " bytes for " << name << ": " << std::get<DataOrigin>(outputSpec).template as<std::string>() << "/" << std::get<DataDescription>(outputSpec).template as<std::string>() << "/" << std::get<2>(outputSpec);
690692
}
691693
std::chrono::time_point<std::chrono::high_resolution_clock> start, end;
692694
if (mVerbosity) {
693695
start = std::chrono::high_resolution_clock::now();
694696
}
695697
buffer.first.emplace(pc.outputs().make<DataAllocator::UninitializedVector<outputDataType>>(std::make_from_tuple<Output>(outputSpec), size));
698+
outputsCreated.insert(name);
696699
if (mVerbosity) {
697700
end = std::chrono::high_resolution_clock::now();
698701
std::chrono::duration<double> elapsed_seconds = end - start;
@@ -705,6 +708,7 @@ void GPURecoWorkflowSpec::run(ProcessingContext& pc)
705708
buffer.first.emplace(pc.outputs().make<DataAllocator::UninitializedVector<outputDataType>>(std::make_from_tuple<Output>(outputSpec), mConfParam->outputBufferSize));
706709
region.ptrBase = (buffer.second = buffer.first->get().data()) + offset;
707710
region.size = buffer.first->get().size() - offset;
711+
outputsCreated.insert(name);
708712
}
709713
}
710714
};
@@ -842,10 +846,6 @@ void GPURecoWorkflowSpec::run(ProcessingContext& pc)
842846
}
843847
}
844848

845-
if (mConfig->configReconstruction.tpc.occupancyMapTimeBins == 0) {
846-
pc.outputs().make<DataAllocator::UninitializedVector<outputDataType>>({gDataOriginTPC, "TPCOCCUPANCYMAP", 0}, 0u);
847-
}
848-
849849
std::unique_ptr<o2::tpc::ClusterNativeAccess> tmpEmptyClNative;
850850
if (createEmptyOutput) {
851851
memset(&ptrs, 0, sizeof(ptrs));
@@ -971,7 +971,10 @@ void GPURecoWorkflowSpec::run(ProcessingContext& pc)
971971
pc.outputs().snapshot({gDataOriginGPU, "ERRORQA", 0}, mErrorQA);
972972
mErrorQA.clear(); // FIXME: This is a race condition once we run multi-threaded!
973973
}
974-
if (mSpecConfig.tpcTriggerHandling && !mSpecConfig.caClusterer) {
974+
if (mSpecConfig.outputSharedClusterMap && !outputsCreated.contains("TPCOCCUPANCYMAP")) {
975+
pc.outputs().make<DataAllocator::UninitializedVector<outputDataType>>({gDataOriginTPC, "TPCOCCUPANCYMAP", 0}, 0u);
976+
}
977+
if (mSpecConfig.tpcTriggerHandling && !outputsCreated.contains("TRIGGERWORDS")) {
975978
pc.outputs().make<DataAllocator::UninitializedVector<outputDataType>>(Output{gDataOriginTPC, "TRIGGERWORDS", 0}, 0u);
976979
}
977980
mTimer->Stop();

0 commit comments

Comments
 (0)