Skip to content

Commit fa02b66

Browse files
ktfdavidrohr
authored andcommitted
DPL: print channel name when tracing messages
1 parent b518797 commit fa02b66

5 files changed

Lines changed: 82 additions & 53 deletions

File tree

Framework/Core/include/Framework/DataRelayer.h

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,24 @@ class DataRelayer
8181
CompletionPolicy::CompletionOp op;
8282
};
8383

84+
enum struct InputType : int {
85+
Invalid = 0,
86+
Data = 1,
87+
SourceInfo = 2,
88+
DomainInfo = 3
89+
};
90+
91+
struct InputInfo {
92+
InputInfo(size_t p, size_t s, InputType t, ChannelIndex i)
93+
: position(p), size(s), type(t), index(i)
94+
{
95+
}
96+
size_t position;
97+
size_t size;
98+
InputType type;
99+
ChannelIndex index;
100+
};
101+
84102
DataRelayer(CompletionPolicy const&,
85103
std::vector<InputRoute> const& routes,
86104
TimesliceIndex&,
@@ -114,6 +132,7 @@ class DataRelayer
114132
/// Notice that we expect that the header is an O2 Header Stack
115133
RelayChoice relay(void const* rawHeader,
116134
std::unique_ptr<fair::mq::Message>* messages,
135+
InputInfo const& info,
117136
size_t nMessages,
118137
size_t nPayloads = 1,
119138
OnDropCallback onDrop = nullptr);

Framework/Core/src/DataProcessingDevice.cxx

Lines changed: 15 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -1765,29 +1765,15 @@ struct WaitBackpressurePolicy {
17651765
/// boilerplate which the user does not need to care about at top level.
17661766
void DataProcessingDevice::handleData(ServiceRegistryRef ref, InputChannelInfo& info)
17671767
{
1768+
using InputInfo = DataRelayer::InputInfo;
1769+
using InputType = DataRelayer::InputType;
1770+
17681771
auto& context = ref.get<DataProcessorContext>();
17691772
// This is the same id as the upper level function, so we get the events
17701773
// associated with the same interval. We will simply use "handle_data" as
17711774
// the category.
17721775
O2_SIGNPOST_ID_FROM_POINTER(cid, device, &info);
17731776

1774-
enum struct InputType : int {
1775-
Invalid = 0,
1776-
Data = 1,
1777-
SourceInfo = 2,
1778-
DomainInfo = 3
1779-
};
1780-
1781-
struct InputInfo {
1782-
InputInfo(size_t p, size_t s, InputType t)
1783-
: position(p), size(s), type(t)
1784-
{
1785-
}
1786-
size_t position;
1787-
size_t size;
1788-
InputType type;
1789-
};
1790-
17911777
// This is how we validate inputs. I.e. we try to enforce the O2 Data model
17921778
// and we do a few stats. We bind parts as a lambda captured variable, rather
17931779
// than an input, because we do not want the outer loop actually be exposed
@@ -1804,8 +1790,8 @@ void DataProcessingDevice::handleData(ServiceRegistryRef ref, InputChannelInfo&
18041790
results.reserve(parts.Size() / 2);
18051791
size_t nTotalPayloads = 0;
18061792

1807-
auto insertInputInfo = [&results, &nTotalPayloads](size_t position, size_t length, InputType type) {
1808-
results.emplace_back(position, length, type);
1793+
auto insertInputInfo = [&results, &nTotalPayloads](size_t position, size_t length, InputType type, ChannelIndex index) {
1794+
results.emplace_back(position, length, type, index);
18091795
if (type != InputType::Invalid && length > 1) {
18101796
nTotalPayloads += length - 1;
18111797
}
@@ -1817,24 +1803,24 @@ void DataProcessingDevice::handleData(ServiceRegistryRef ref, InputChannelInfo&
18171803
if (sih) {
18181804
O2_SIGNPOST_EVENT_EMIT(device, cid, "handle_data", "Got SourceInfoHeader with state %d", (int)sih->state);
18191805
info.state = sih->state;
1820-
insertInputInfo(pi, 2, InputType::SourceInfo);
1806+
insertInputInfo(pi, 2, InputType::SourceInfo, info.id);
18211807
*context.wasActive = true;
18221808
continue;
18231809
}
18241810
auto dih = o2::header::get<DomainInfoHeader*>(headerData);
18251811
if (dih) {
1826-
insertInputInfo(pi, 2, InputType::DomainInfo);
1812+
insertInputInfo(pi, 2, InputType::DomainInfo, info.id);
18271813
*context.wasActive = true;
18281814
continue;
18291815
}
18301816
auto dh = o2::header::get<DataHeader*>(headerData);
18311817
if (!dh) {
1832-
insertInputInfo(pi, 0, InputType::Invalid);
1818+
insertInputInfo(pi, 0, InputType::Invalid, info.id);
18331819
O2_SIGNPOST_EVENT_EMIT_ERROR(device, cid, "handle_data", "Header is not a DataHeader?");
18341820
continue;
18351821
}
18361822
if (dh->payloadSize > parts.At(pi + 1)->GetSize()) {
1837-
insertInputInfo(pi, 0, InputType::Invalid);
1823+
insertInputInfo(pi, 0, InputType::Invalid, info.id);
18381824
O2_SIGNPOST_EVENT_EMIT_ERROR(device, cid, "handle_data", "DataHeader payloadSize mismatch");
18391825
continue;
18401826
}
@@ -1847,14 +1833,14 @@ void DataProcessingDevice::handleData(ServiceRegistryRef ref, InputChannelInfo&
18471833
O2_SIGNPOST_START(parts, pid, "parts", "Processing DataHeader with splitPayloadParts %d and splitPayloadIndex %d", dh->splitPayloadParts, dh->splitPayloadIndex);
18481834
}
18491835
if (!dph) {
1850-
insertInputInfo(pi, 2, InputType::Invalid);
1836+
insertInputInfo(pi, 2, InputType::Invalid, info.id);
18511837
O2_SIGNPOST_EVENT_EMIT_ERROR(device, cid, "handle_data", "Header stack does not contain DataProcessingHeader");
18521838
continue;
18531839
}
18541840
if (dh->splitPayloadParts > 0 && dh->splitPayloadParts == dh->splitPayloadIndex) {
18551841
// this is indicating a sequence of payloads following the header
18561842
// FIXME: we will probably also set the DataHeader version
1857-
insertInputInfo(pi, dh->splitPayloadParts + 1, InputType::Data);
1843+
insertInputInfo(pi, dh->splitPayloadParts + 1, InputType::Data, info.id);
18581844
pi += dh->splitPayloadParts - 1;
18591845
} else {
18601846
// We can set the type for the next splitPayloadParts
@@ -1864,12 +1850,12 @@ void DataProcessingDevice::handleData(ServiceRegistryRef ref, InputChannelInfo&
18641850
size_t finalSplitPayloadIndex = pi + (dh->splitPayloadParts > 0 ? dh->splitPayloadParts : 1) * 2;
18651851
if (finalSplitPayloadIndex > parts.Size()) {
18661852
O2_SIGNPOST_EVENT_EMIT_ERROR(device, cid, "handle_data", "DataHeader::splitPayloadParts invalid");
1867-
insertInputInfo(pi, 0, InputType::Invalid);
1853+
insertInputInfo(pi, 0, InputType::Invalid, info.id);
18681854
continue;
18691855
}
1870-
insertInputInfo(pi, 2, InputType::Data);
1856+
insertInputInfo(pi, 2, InputType::Data, info.id);
18711857
for (; pi + 2 < finalSplitPayloadIndex; pi += 2) {
1872-
insertInputInfo(pi + 2, 2, InputType::Data);
1858+
insertInputInfo(pi + 2, 2, InputType::Data, info.id);
18731859
}
18741860
}
18751861
}
@@ -1941,6 +1927,7 @@ void DataProcessingDevice::handleData(ServiceRegistryRef ref, InputChannelInfo&
19411927
};
19421928
auto relayed = relayer.relay(parts.At(headerIndex)->GetData(),
19431929
&parts.At(headerIndex),
1930+
input,
19441931
nMessages,
19451932
nPayloadsPerHeader,
19461933
onDrop);

Framework/Core/src/DataRelayer.cxx

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -394,6 +394,7 @@ void DataRelayer::pruneCache(TimesliceSlot slot, OnDropCallback onDrop)
394394
DataRelayer::RelayChoice
395395
DataRelayer::relay(void const* rawHeader,
396396
std::unique_ptr<fair::mq::Message>* messages,
397+
InputInfo const& info,
397398
size_t nMessages,
398399
size_t nPayloads,
399400
std::function<void(TimesliceSlot, std::vector<MessageSet>&, TimesliceIndex::OldestOutputInfo)> onDrop)
@@ -443,11 +444,13 @@ DataRelayer::RelayChoice
443444
&nMessages,
444445
&nPayloads,
445446
&cache = mCache,
446-
numInputTypes = mDistinctRoutesIndex.size()](TimesliceId timeslice, int input, TimesliceSlot slot) {
447+
&services = mContext,
448+
numInputTypes = mDistinctRoutesIndex.size()](TimesliceId timeslice, int input, TimesliceSlot slot, InputInfo const& info) {
447449
O2_SIGNPOST_ID_GENERATE(aid, data_relayer);
448-
O2_SIGNPOST_EVENT_EMIT(data_relayer, aid, "saveInSlot", "saving %{public}s@%zu in slot %zu",
450+
O2_SIGNPOST_EVENT_EMIT(data_relayer, aid, "saveInSlot", "saving %{public}s@%zu in slot %zu from %{public}s",
449451
fmt::format("{:x}", *o2::header::get<DataHeader*>(messages[0]->GetData())).c_str(),
450-
timeslice.value, slot.index);
452+
timeslice.value, slot.index,
453+
info.index.value == ChannelIndex::INVALID ? "invalid" : services.get<FairMQDeviceProxy>().getInputChannel(info.index)->GetName().c_str());
451454
auto cacheIdx = numInputTypes * slot.index + input;
452455
MessageSet& target = cache[cacheIdx];
453456
cachedStateMetrics[cacheIdx] = CacheEntryStatus::PENDING;
@@ -537,7 +540,7 @@ DataRelayer::RelayChoice
537540
this->pruneCache(slot, onDrop);
538541
mPruneOps.erase(std::remove_if(mPruneOps.begin(), mPruneOps.end(), [slot](const auto& x) { return x.slot == slot; }), mPruneOps.end());
539542
}
540-
saveInSlot(timeslice, input, slot);
543+
saveInSlot(timeslice, input, slot, info);
541544
index.publishSlot(slot);
542545
index.markAsDirty(slot, true);
543546
stats.updateStats({static_cast<short>(ProcessingStatsId::RELAYED_MESSAGES), DataProcessingStats::Op::Add, (int)1});
@@ -619,7 +622,7 @@ DataRelayer::RelayChoice
619622
// cache still holds the old data, so we prune it.
620623
this->pruneCache(slot, onDrop);
621624
mPruneOps.erase(std::remove_if(mPruneOps.begin(), mPruneOps.end(), [slot](const auto& x) { return x.slot == slot; }), mPruneOps.end());
622-
saveInSlot(timeslice, input, slot);
625+
saveInSlot(timeslice, input, slot, info);
623626
index.publishSlot(slot);
624627
index.markAsDirty(slot, true);
625628
return RelayChoice{.type = RelayChoice::Type::WillRelay};

Framework/Core/test/benchmark_DataRelayer.cxx

Lines changed: 12 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -86,8 +86,9 @@ static void BM_RelaySingleSlot(benchmark::State& state)
8686
inflightMessages.emplace_back(transport->CreateMessage(1000));
8787
memcpy(inflightMessages[0]->GetData(), stack.data(), stack.size());
8888

89+
DataRelayer::InputInfo fakeInfo{0, inflightMessages.size(), DataRelayer::InputType::Data, {ChannelIndex::INVALID}};
8990
for (auto _ : state) {
90-
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), inflightMessages.size());
91+
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size());
9192
std::vector<RecordAction> ready;
9293
relayer.getReadyToProcess(ready);
9394
assert(ready.size() == 1);
@@ -144,7 +145,8 @@ static void BM_RelayMultipleSlots(benchmark::State& state)
144145
Stack stack{dh, DataProcessingHeader{timeslice++, 1}};
145146
memcpy(inflightMessages[0]->GetData(), stack.data(), stack.size());
146147

147-
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), inflightMessages.size());
148+
DataRelayer::InputInfo fakeInfo{0, inflightMessages.size(), DataRelayer::InputType::Data, {ChannelIndex::INVALID}};
149+
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size());
148150
std::vector<RecordAction> ready;
149151
relayer.getReadyToProcess(ready);
150152
assert(ready.size() == 1);
@@ -211,13 +213,15 @@ static void BM_RelayMultipleRoutes(benchmark::State& state)
211213
memcpy(inflightMessages[2]->GetData(), stack2.data(), stack2.size());
212214

213215
for (auto _ : state) {
214-
relayer.relay(inflightMessages[0]->GetData(), &inflightMessages[0], 2);
216+
DataRelayer::InputInfo fakeInfo{0, inflightMessages.size(), DataRelayer::InputType::Data, {ChannelIndex::INVALID}};
217+
relayer.relay(inflightMessages[0]->GetData(), &inflightMessages[0], fakeInfo, 2);
215218
std::vector<RecordAction> ready;
216219
relayer.getReadyToProcess(ready);
217220
assert(ready.size() == 1);
218221
assert(ready[0].op == CompletionPolicy::CompletionOp::Consume);
219222

220-
relayer.relay(inflightMessages[2]->GetData(), &inflightMessages[2], 2);
223+
DataRelayer::InputInfo fakeInfo2{0, inflightMessages.size(), DataRelayer::InputType::Data, {ChannelIndex::INVALID}};
224+
relayer.relay(inflightMessages[2]->GetData(), &inflightMessages[2], fakeInfo2, 2);
221225
ready.clear();
222226
relayer.getReadyToProcess(ready);
223227
assert(ready.size() == 1);
@@ -282,8 +286,9 @@ static void BM_RelaySplitParts(benchmark::State& state)
282286
inflightMessages.emplace_back(std::move(payload));
283287
}
284288

289+
DataRelayer::InputInfo fakeInfo{0, inflightMessages.size(), DataRelayer::InputType::Data, {ChannelIndex::INVALID}};
285290
for (auto _ : state) {
286-
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), inflightMessages.size());
291+
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size());
287292
std::vector<RecordAction> ready;
288293
relayer.getReadyToProcess(ready);
289294
assert(ready.size() == 1);
@@ -336,8 +341,9 @@ static void BM_RelayMultiplePayloads(benchmark::State& state)
336341
inflightMessages.emplace_back(transport->CreateMessage(dh.payloadSize));
337342
}
338343

344+
DataRelayer::InputInfo fakeInfo{0, inflightMessages.size(), DataRelayer::InputType::Data, {ChannelIndex::INVALID}};
339345
for (auto _ : state) {
340-
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), inflightMessages.size(), nPayloads);
346+
relayer.relay(inflightMessages[0]->GetData(), inflightMessages.data(), fakeInfo, inflightMessages.size(), nPayloads);
341347
std::vector<RecordAction> ready;
342348
relayer.getReadyToProcess(ready);
343349
assert(ready.size() == 1);

0 commit comments

Comments
 (0)