Skip to content

Commit 7c80e4f

Browse files
committed
DPL: Add DPL_OUTPUT_PROXY_WHENANY env variable setting for output proxy
1 parent f0c43c4 commit 7c80e4f

3 files changed

Lines changed: 13 additions & 0 deletions

File tree

Framework/Core/include/Framework/CompletionPolicyHelpers.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,8 @@ struct CompletionPolicyHelpers {
4949
{
5050
return consumeWhenAny("consume-any", matcher);
5151
}
52+
static CompletionPolicy consumeWhenAny(std::string matchName);
53+
5254
/// When any of the parts of the record have been received, process the existing and free the associated payloads.
5355
/// This allows freeing things as early as possible, while still being able to wait
5456
/// all the parts before disposing the timeslice completely

Framework/Core/src/CompletionPolicyHelpers.cxx

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -204,6 +204,14 @@ CompletionPolicy CompletionPolicyHelpers::consumeWhenAny(const char* name, Compl
204204
return CompletionPolicy{name, matcher, callback};
205205
}
206206

207+
CompletionPolicy CompletionPolicyHelpers::consumeWhenAny(std::string matchName)
208+
{
209+
auto matcher = [matchName](DeviceSpec const& device) -> bool {
210+
return std::regex_match(device.name.begin(), device.name.end(), std::regex(matchName));
211+
};
212+
return consumeWhenAny(matcher);
213+
}
214+
207215
CompletionPolicy CompletionPolicyHelpers::processWhenAny(const char* name, CompletionPolicy::Matcher matcher)
208216
{
209217
auto callback = [](InputSpan const& inputs) -> CompletionPolicy::CompletionOp {

Framework/Utils/src/dpl-output-proxy.cxx

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,8 +59,11 @@ void customize(std::vector<ConfigParamSpec>& workflowOptions)
5959
void customize(std::vector<o2::framework::CompletionPolicy>& policies)
6060
{
6161
static bool doOrdered = getenv("DPL_OUTPUT_PROXY_ORDERED") && atoi(getenv("DPL_OUTPUT_PROXY_ORDERED"));
62+
static bool doAny = getenv("DPL_OUTPUT_PROXY_WHENANY") && atoi(getenv("DPL_OUTPUT_PROXY_WHENANY"));
6263
if (doOrdered) {
6364
policies.push_back(CompletionPolicyHelpers::consumeWhenAllOrdered(".*"));
65+
} else if (doAny) {
66+
policies.push_back(CompletionPolicyHelpers::consumeWhenAny(".*"));
6467
}
6568
}
6669

0 commit comments

Comments
 (0)