Skip to content

Commit db201de

Browse files
committed
Clean up remaining issues for TF rate throttling
1 parent be7269a commit db201de

7 files changed

Lines changed: 37 additions & 12 deletions

File tree

Framework/Core/include/Framework/CommonDataProcessors.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,8 @@ struct CommonDataProcessors {
6060

6161
/// @return a dummy DataProcessorSpec which requires all the passed @a InputSpec
6262
/// and simply discards them.
63-
static DataProcessorSpec getDummySink(std::vector<InputSpec> const& danglingInputs);
63+
static DataProcessorSpec getDummySink(std::vector<InputSpec> const& danglingInputs,
64+
int rateLimitingIPCID);
6465
};
6566

6667
} // namespace o2::framework

Framework/Core/src/CommonDataProcessors.cxx

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -501,7 +501,7 @@ DataProcessorSpec CommonDataProcessors::getGlobalFairMQSink(std::vector<InputSpe
501501
return specifyFairMQDeviceOutputProxy("internal-dpl-injected-output-proxy", danglingOutputInputs, defaultChannelConfig.c_str());
502502
}
503503

504-
DataProcessorSpec CommonDataProcessors::getDummySink(std::vector<InputSpec> const& danglingOutputInputs)
504+
DataProcessorSpec CommonDataProcessors::getDummySink(std::vector<InputSpec> const& danglingOutputInputs, int rateLimitingIPCID)
505505
{
506506
return DataProcessorSpec{
507507
.name = "internal-dpl-injected-dummy-sink",
@@ -523,7 +523,13 @@ DataProcessorSpec CommonDataProcessors::getDummySink(std::vector<InputSpec> cons
523523

524524
return adaptStateless([]() {
525525
});
526-
})}};
526+
})},
527+
.options = rateLimitingIPCID != -1 ? std::vector<ConfigParamSpec>{{"channel-config", VariantType::String, // raw input channel
528+
"name=metric-feedback,type=push,method=bind,address=ipc://@metric-feedback-" + std::to_string(rateLimitingIPCID) + ",transport=shmem,rateLogging=10",
529+
{"Out-of-band channel config"}}}
530+
: std::vector<ConfigParamSpec>()
531+
532+
};
527533
}
528534

529535
#pragma GCC diagnostic pop

Framework/Core/src/DataProcessingDevice.cxx

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -536,7 +536,7 @@ void DataProcessingDevice::fillContext(DataProcessorContext& context, DeviceCont
536536
deviceContext.stats = &mStats;
537537
context.isSink = false;
538538
// If nothing is a sink, the rate limiting simply does not trigger.
539-
int enableRateLimiting = std::stoi(fConfig->GetValue<std::string>("timeframes-rate-limit"));
539+
bool enableRateLimiting = std::stoi(fConfig->GetValue<std::string>("timeframes-rate-limit"));
540540
if (enableRateLimiting) {
541541
for (auto& spec : mSpec.outputs) {
542542
if (spec.matcher.binding.value == "dpl-summary") {

Framework/Core/src/ExternalFairMQDeviceProxy.cxx

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -415,17 +415,29 @@ DataProcessorSpec specifyExternalFairMQDeviceProxy(char const* name,
415415
static int64_t consumedTimeframes = 0;
416416
static int64_t sentTimeframes = 0;
417417
auto device = ctx.services().get<RawDeviceService>().device();
418-
if (device->fChannels.count("metric-feedback")) {
418+
int maxTFInFlight = std::stoi(device->fConfig->GetValue<std::string>("timeframes-rate-limit"));
419+
if (maxTFInFlight && device->fChannels.count("metric-feedback")) {
419420
// FIXME: only two channels difference
420-
while ((sentTimeframes - consumedTimeframes) > 2) {
421-
FairMQMessagePtr msg;
422-
auto count = device->Receive(msg, "metric-feedback", 0, -1);
421+
int waitMessage = 0;
422+
int recvTimeot = 0;
423+
while ((sentTimeframes - consumedTimeframes) >= maxTFInFlight) {
424+
if (recvTimeot == -1 && waitMessage == 0) {
425+
LOG(alarm) << "Maximum number of TF in flight reached (" << maxTFInFlight << ": published " << sentTimeframes << " - finished " << consumedTimeframes << "), waiting";
426+
waitMessage = 1;
427+
}
428+
auto msg = device->NewMessageFor("metric-feedback", 0, 0);
429+
430+
auto count = device->Receive(msg, "metric-feedback", 0, recvTimeot);
423431
if (count <= 0) {
424-
return;
432+
recvTimeot = -1;
433+
continue;
425434
}
426-
assert(msg->GetSize() == 8);
435+
assert(msg->GetSize() == 8);
427436
consumedTimeframes = *(int64_t*)msg->GetData();
428437
}
438+
if (waitMessage) {
439+
LOG(important) << (sentTimeframes - consumedTimeframes) << " / " << maxTFInFlight << " TF in flight, continuing to publish";
440+
}
429441
}
430442

431443
FairMQParts parts;

Framework/Core/src/WorkflowCustomizationHelpers.cxx

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,9 @@ std::vector<ConfigParamSpec> WorkflowCustomizationHelpers::requiredWorkflowOptio
4343
ConfigParamSpec{"labels", VariantType::String, "", {"add labels to dataprocessors"}},
4444
ConfigParamSpec{"workflow-suffix", VariantType::String, "", {"suffix to add to all dataprocessors"}},
4545

46+
// options for TF rate limiting
47+
ConfigParamSpec{"timeframes-rate-limit-ipcid", VariantType::String, "-1", {"Suffix for IPC channel for metrix-feedback, -1 = disable"}},
48+
4649
// options for AOD rate limiting
4750
ConfigParamSpec{"aod-memory-rate-limit", VariantType::Int64, 0LL, {"Rate limit AOD processing based on memory"}},
4851

Framework/Core/src/WorkflowHelpers.cxx

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -567,7 +567,8 @@ void WorkflowHelpers::injectServiceDevices(WorkflowSpec& workflow, ConfigContext
567567
if (unmatched.size() > 0 || redirectedOutputsInputs.size() > 0) {
568568
std::vector<InputSpec> ignored = unmatched;
569569
ignored.insert(ignored.end(), redirectedOutputsInputs.begin(), redirectedOutputsInputs.end());
570-
extraSpecs.push_back(CommonDataProcessors::getDummySink(ignored));
570+
int rateLimitingIPCID = std::stoi(ctx.options().get<std::string>("timeframes-rate-limit-ipcid"));
571+
extraSpecs.push_back(CommonDataProcessors::getDummySink(ignored, rateLimitingIPCID));
571572
}
572573

573574
workflow.insert(workflow.end(), extraSpecs.begin(), extraSpecs.end());

GPU/GPUTracking/DataTypes/GPUMemorySizeScalers.cxx

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,9 @@ void GPUMemorySizeScalers::rescaleMaxMem(size_t newAvailableMemory)
2121
{
2222
GPUMemorySizeScalers tmp;
2323
double scaleFactor = (double)newAvailableMemory / tmp.availableMemory;
24-
GPUInfo("Rescaling buffer size limits from %lu to %lu bytes of memory (factor %f)", tmp.availableMemory, newAvailableMemory, scaleFactor);
24+
if (scaleFactor != 1.) {
25+
GPUInfo("Rescaling buffer size limits from %lu to %lu bytes of memory (factor %f)", tmp.availableMemory, newAvailableMemory, scaleFactor);
26+
}
2527
tpcMaxPeaks = (double)tmp.tpcMaxPeaks * scaleFactor;
2628
tpcMaxClusters = (double)tmp.tpcMaxClusters * scaleFactor;
2729
tpcMaxStartHits = (double)tmp.tpcMaxStartHits * scaleFactor;

0 commit comments

Comments
 (0)