Skip to content

Commit 5288e8f

Browse files
committed
ensure empty frame is indeed empty
1 parent 12917b5 commit 5288e8f

7 files changed

Lines changed: 135 additions & 88 deletions

File tree

DataFormats/Detectors/TRD/include/DataFormatsTRD/EventRecord.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
#include "CommonDataFormat/RangeReference.h"
1919
#include "FairLogger.h"
2020
#include "DataFormatsTRD/Tracklet64.h"
21+
#include "Framework/ProcessingContext.h"
2122

2223
namespace o2::trd
2324
{
@@ -93,6 +94,9 @@ class EventStorage
9394
void addTracklets(InteractionRecord& ir, std::vector<Tracklet64>& tracklets);
9495
void addTracklets(InteractionRecord& ir, std::vector<Tracklet64>::iterator& start, std::vector<Tracklet64>::iterator& end);
9596
void unpackDataForSending(std::vector<TriggerRecord>& triggers, std::vector<Tracklet64>& tracklets, std::vector<Digit>& digits);
97+
void sendData(o2::framework::ProcessingContext& pc);
98+
//this could replace by keeing a running total on addition TODO
99+
void sumTrackletsDigitsTriggers(uint64_t& tracklets, uint64_t& digits, uint64_t& triggers);
96100
int sumTracklets();
97101
int sumDigits();
98102
std::vector<Tracklet64>& getTracklets(InteractionRecord& ir);

DataFormats/Detectors/TRD/src/EventRecord.cxx

Lines changed: 54 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,18 @@
2121
#include "DataFormatsTRD/Digit.h"
2222
#include "DataFormatsTRD/EventRecord.h"
2323
#include "DataFormatsTRD/Constants.h"
24+
25+
#include "Framework/Output.h"
26+
#include "Framework/ProcessingContext.h"
27+
#include "Framework/ControlService.h"
28+
#include "Framework/ConfigParamRegistry.h"
29+
#include "Framework/RawDeviceService.h"
30+
#include "Framework/DeviceSpec.h"
31+
#include "Framework/DataSpecUtils.h"
32+
#include "Framework/InputRecordWalker.h"
33+
34+
#include "DataFormatsTRD/Constants.h"
35+
2436
#include <cassert>
2537
#include <array>
2638
#include <string>
@@ -142,14 +154,43 @@ void EventStorage::unpackDataForSending(std::vector<TriggerRecord>& triggers, st
142154
{
143155
int digitcount = 0;
144156
int trackletcount = 0;
145-
for (auto event : mEventRecords) {
157+
for (auto& event : mEventRecords) {
158+
tracklets.insert(std::end(tracklets), std::begin(event.getTracklets()), std::end(event.getTracklets()));
159+
digits.insert(std::end(digits), std::begin(event.getDigits()), std::end(event.getDigits()));
160+
triggers.emplace_back(event.getBCData(), digitcount, event.getDigits().size(), trackletcount, event.getTracklets().size());
161+
digitcount += event.getDigits().size();
162+
trackletcount += event.getTracklets().size();
163+
}
164+
}
165+
166+
void EventStorage::sendData(o2::framework::ProcessingContext& pc)
167+
{
168+
//at this point we know the total number of tracklets and digits and triggers.
169+
//hence we can create the relevant objects inside the message as opposed to creating a local object and snapshotting it(copying) it
170+
//into the message.
171+
uint64_t trackletcount = 0;
172+
uint64_t digitcount = 0;
173+
uint64_t triggercount = 0;
174+
sumTrackletsDigitsTriggers(trackletcount, digitcount, triggercount);
175+
std::vector<Tracklet64> tracklets;
176+
tracklets.reserve(trackletcount);
177+
std::vector<Digit> digits;
178+
digits.reserve(digitcount);
179+
std::vector<TriggerRecord> triggers;
180+
triggers.reserve(triggercount);
181+
for (auto& event : mEventRecords) {
146182
tracklets.insert(std::end(tracklets), std::begin(event.getTracklets()), std::end(event.getTracklets()));
147183
digits.insert(std::end(digits), std::begin(event.getDigits()), std::end(event.getDigits()));
148184
triggers.emplace_back(event.getBCData(), digitcount, event.getDigits().size(), trackletcount, event.getTracklets().size());
149185
digitcount += event.getDigits().size();
150186
trackletcount += event.getTracklets().size();
151187
}
188+
LOG(info) << "Sending data onwards with " << digits.size() << " Digits and " << tracklets.size() << " Tracklets and " << triggers.size() << " Triggers";
189+
pc.outputs().snapshot(o2::framework::Output{o2::header::gDataOriginTRD, "DIGITS", 0, o2::framework::Lifetime::Timeframe}, digits);
190+
pc.outputs().snapshot(o2::framework::Output{o2::header::gDataOriginTRD, "TRACKLETS", 0, o2::framework::Lifetime::Timeframe}, tracklets);
191+
pc.outputs().snapshot(o2::framework::Output{o2::header::gDataOriginTRD, "TRKTRGRD", 0, o2::framework::Lifetime::Timeframe}, triggers);
152192
}
193+
153194
int EventStorage::sumTracklets()
154195
{
155196
int sum = 0;
@@ -166,6 +207,18 @@ int EventStorage::sumDigits()
166207
}
167208
return sum;
168209
}
210+
void EventStorage::sumTrackletsDigitsTriggers(uint64_t& tracklets, uint64_t& digits, uint64_t& triggers)
211+
{
212+
int digitsum = 0;
213+
int trackletsum = 0;
214+
int triggersum = 0;
215+
for (auto event : mEventRecords) {
216+
digitsum += event.getDigits().size();
217+
trackletsum += event.getTracklets().size();
218+
triggersum++;
219+
}
220+
}
221+
169222
std::vector<Tracklet64>& EventStorage::getTracklets(InteractionRecord& ir)
170223
{
171224
bool found = false;

Detectors/TRD/reconstruction/include/TRDReconstruction/CruRawReader.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,8 @@ class CruRawReader
9292
std::vector<Digit>& getDigits(InteractionRecord& ir) { return mEventRecords.getDigits(ir); };
9393
// std::vector<o2::trd::TriggerRecord> getIR() { return mEventTriggers; }
9494
void getParsedObjects(std::vector<Tracklet64>& tracklets, std::vector<Digit>& cdigits, std::vector<TriggerRecord>& triggers);
95+
void getParsedObjectsandClear(std::vector<Tracklet64>& tracklets, std::vector<Digit>& digits, std::vector<TriggerRecord>& triggers);
96+
void buildDPLOutputs(o2::framework::ProcessingContext& outputs);
9597
int getDigitsFound() { return mTotalDigitsFound; }
9698
int getTrackletsFound() { return mTotalTrackletsFound; }
9799
int sumTrackletsFound() { return mEventRecords.sumTracklets(); }

Detectors/TRD/reconstruction/include/TRDReconstruction/DataReaderTask.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,8 @@ class DataReaderTask : public Task
5757
bool mHeaderVerbose{false}; // verbose output of headers
5858
bool mCompressedData{false}; // are we dealing with the compressed data from the flp (send via option)
5959
bool mByteSwap{true}; // whether we are to byteswap the incoming data, mc is not byteswapped, raw data is (too be changed in cru at some point)
60-
o2::header::DataDescription mDataSpec; // input spec of th raw incoming data
60+
// o2::header::DataDescription mDataDesc; // Data description of the incoming data
61+
std::string mDataDesc;
6162
};
6263

6364
} // namespace o2::trd

Detectors/TRD/reconstruction/src/CruRawReader.cxx

Lines changed: 23 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,14 @@
2323
#include "TRDReconstruction/TrackletsParser.h"
2424
#include "DataFormatsTRD/Constants.h"
2525

26+
#include "Framework/ControlService.h"
27+
#include "Framework/ConfigParamRegistry.h"
28+
#include "Framework/RawDeviceService.h"
29+
#include "Framework/DeviceSpec.h"
30+
#include "Framework/DataSpecUtils.h"
31+
#include "Framework/Output.h"
32+
#include "Framework/InputRecordWalker.h"
33+
2634
#include <cstring>
2735
#include <string>
2836
#include <vector>
@@ -415,31 +423,20 @@ void CruRawReader::getParsedObjects(std::vector<Tracklet64>& tracklets, std::vec
415423
int digitcountsum = 0;
416424
int trackletcountsum = 0;
417425
mEventRecords.unpackDataForSending(triggers, tracklets, digits);
418-
/*for(auto eventrecord: mEventRecords)//loop over triggers incase they have already been done.
419-
{
420-
int digitcount=0;
421-
int trackletcount=0;
422-
int start,end;
423-
LOG(info) << __func__ << " " << tracklets.size() << " "
424-
<< cdigits.size()<< " trackletv size:"<< mEventRecords.getTracklets(ir.getBCData());
425-
for(auto trackletv: mEventStores.getTracklets(ir.getBCData())){
426-
//loop through the vector of ranges
427-
start=trackletv.getFirstEntry();
428-
end= start+trackletv.getEntries();
429-
LOG(info) << "insert tracklets from " << start<< " " << end;
430-
tracklets.insert(tracklets.end(),mEventTracklets.begin()+start, mEventTracklets.begin()+end);
431-
trackletcount+=trackletv.getEntries();
432-
}
433-
for(auto digitv: mEventStores.getDigits(ir.getBCData())){
434-
start=digitv.getFirstEntry();
435-
end= start+digitv.getEntries();
436-
LOG(info) << "insert digits from " << start<< " " << end;
437-
cdigits.insert(cdigits.end(),mEventCompressedDigits.begin()+start , mEventCompressedDigits.begin()+end);
438-
digitcount+=digitv.getEntries();
439-
}
440-
triggers.emplace_back(ir.getBCData(),digitcountsum,digitcount,trackletcountsum,trackletcount);
441-
digitcountsum+=digitcount;
442-
trackletcountsum+=trackletcount;
443-
}*/
444426
}
427+
428+
void CruRawReader::getParsedObjectsandClear(std::vector<Tracklet64>& tracklets, std::vector<Digit>& digits, std::vector<TriggerRecord>& triggers)
429+
{
430+
getParsedObjects(tracklets, digits, triggers);
431+
clearall();
432+
}
433+
434+
//write the output data directly to the given DataAllocator from the datareader task.
435+
void CruRawReader::buildDPLOutputs(o2::framework::ProcessingContext& pc)
436+
{
437+
mEventRecords.sendData(pc);
438+
// pc.outputs().snapshot(Output{o2::header::gDataOriginTRD,"STATS",0,Lifetime::Timerframe},mStats);
439+
clearall(); // having now written the messages clear for next.
440+
}
441+
445442
} // namespace o2::trd

Detectors/TRD/reconstruction/src/DataReader.cxx

Lines changed: 1 addition & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -76,23 +76,11 @@ WorkflowSpec defineDataProcessing(ConfigContext const& cfgc)
7676

7777
WorkflowSpec workflow;
7878

79-
/*
80-
* This is originally replicated from TOF
81-
We define at run time the number of devices to be attached
82-
to the workflow and the input matching string of the device.
83-
This is is done with a configuration string like the following
84-
one, where the input matching for each device is provide in
85-
comma-separated strings. For instance
86-
*/
87-
88-
// std::stringstream ssconfig(inputspec);
8979
std::string iconfig;
9080
std::string inputDescription;
9181
int idevice = 0;
92-
// LOG(info) << "expected incoming data definition : " << inputspec;
93-
// this is probably never going to be used but would to nice to know hence here.
9482
auto orig = o2::header::gDataOriginTRD;
95-
auto inputs = o2::framework::select(inputspec.c_str());
83+
auto inputs = o2::framework::select(std::string("x:TRD/" + inputspec).c_str());
9684
for (auto& inp : inputs) {
9785
// take care of case where our data is not in the time frame
9886
inp.lifetime = Lifetime::Optional;

Detectors/TRD/reconstruction/src/DataReaderTask.cxx

Lines changed: 49 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -47,14 +47,15 @@ void DataReaderTask::sendData(ProcessingContext& pc, bool blankframe)
4747
{
4848
// mReader.getParsedObjects(mTracklets,mDigits,mTriggers);
4949
if (!blankframe) {
50-
mReader.getParsedObjects(mTracklets, mDigits, mTriggers);
50+
mReader.buildDPLOutputs(pc); //getParsedObjectsandClear(mTracklets, mDigits, mTriggers);
51+
} else {
52+
//ensure the objects we are sending back are indeed blank.
53+
LOG(info) << "Sending data onwards with " << mDigits.size() << " Digits and " << mTracklets.size() << " Tracklets and " << mTriggers.size() << " Triggers and blankframe:" << blankframe;
54+
pc.outputs().snapshot(Output{o2::header::gDataOriginTRD, "DIGITS", 0, Lifetime::Timeframe}, mDigits);
55+
pc.outputs().snapshot(Output{o2::header::gDataOriginTRD, "TRACKLETS", 0, Lifetime::Timeframe}, mTracklets);
56+
pc.outputs().snapshot(Output{o2::header::gDataOriginTRD, "TRKTRGRD", 0, Lifetime::Timeframe}, mTriggers);
57+
// pc.outputs().snapshot(Output{o2::header::gDataOriginTRD,"STATS",0,Lifetime::Timerframe},mStats);
5158
}
52-
53-
LOG(info) << "Sending data onwards with " << mDigits.size() << " Digits and " << mTracklets.size() << " Tracklets and " << mTriggers.size() << " Triggers and blankframe:" << blankframe;
54-
pc.outputs().snapshot(Output{o2::header::gDataOriginTRD, "DIGITS", 0, Lifetime::Timeframe}, mDigits);
55-
pc.outputs().snapshot(Output{o2::header::gDataOriginTRD, "TRACKLETS", 0, Lifetime::Timeframe}, mTracklets);
56-
pc.outputs().snapshot(Output{o2::header::gDataOriginTRD, "TRKTRGRD", 0, Lifetime::Timeframe}, mTriggers);
57-
// pc.outputs().snapshot(Output{o2::header::gDataOriginTRD,"STATS",0,Lifetime::Timerframe},mStats);
5859
}
5960

6061
void DataReaderTask::run(ProcessingContext& pc)
@@ -67,16 +68,16 @@ void DataReaderTask::run(ProcessingContext& pc)
6768
auto device = pc.services().get<o2::framework::RawDeviceService>().device();
6869
auto outputRoutes = pc.services().get<o2::framework::RawDeviceService>().spec().outputs;
6970
auto fairMQChannel = outputRoutes.at(0).channel;
70-
mDataSpec = o2::header::gDataDescriptionRawData;
7171

72-
std::vector<InputSpec> dummy{InputSpec{"dummy", ConcreteDataMatcher{"TRD", mDataSpec, 0xDEADBEEF}}};
73-
// if we see requested data type input with 0xDEADBEEF subspec and 0 payload this means that the "delayed message"
72+
std::vector<InputSpec> dummy{InputSpec{"dummy", ConcreteDataMatcher{"TRD", "RAWDATA", 0xDEADBEEF}}};
73+
//std::vector<InputSpec> dummy{InputSpec{"dummy", ConcreteDataMatcher{"TRD","RAWDATA"/* mDataDesc.c_str()*/, 0xDEADBEEF}}};
74+
// if we see requested data type input with 0xDEADBEEF subspec and 0 payload this mecans that the "delayed message"
7475
// // mechanism created it in absence of real data from upstream. Processor should send empty output to not block the workflow
7576

7677
for (const auto& ref : InputRecordWalker(pc.inputs(), dummy)) {
7778
const auto dh = o2::framework::DataRefUtils::getHeader<o2::header::DataHeader*>(ref);
78-
if (dh->payloadSize == 16 || dh->payloadSize == 0) {
79-
LOGP(WARNING, "Found input [{}/{}/{:#x}] TF#{} 1st_orbit:{} Payload {} : assuming no payload for all links in this TF",
79+
if (dh->payloadSize == 0) { //}|| dh->payloadSize==16) {
80+
LOGP(WARNING, "Found blank input input [{}/{}/{:#x}] TF#{} 1st_orbit:{} Payload {} : ",
8081
dh->dataOrigin.str, dh->dataDescription.str, dh->subSpecification, dh->tfCounter, dh->firstTForbit, dh->payloadSize);
8182
sendData(pc, true); //send the empty tf data.
8283
return;
@@ -94,49 +95,50 @@ void DataReaderTask::run(ProcessingContext& pc)
9495
for (auto const& ref : iit) {
9596
if (mVerbose) {
9697
const auto dh = DataRefUtils::getHeader<o2::header::DataHeader*>(ref);
97-
LOGP(info, "Found input [{}/{}/{:#x}] TF#{} 1st_orbit:{} Payload {} : assuming no payload for all links in this TF",
98+
LOGP(info, "Found input [{}/{}/{:#x}] TF#{} 1st_orbit:{} Payload {} : ",
9899
dh->dataOrigin.str, dh->dataDescription.str, dh->subSpecification, dh->tfCounter, dh->firstTForbit, dh->payloadSize);
99100
}
100101
const auto* headerIn = DataRefUtils::getHeader<o2::header::DataHeader*>(ref);
101102
auto payloadIn = ref.payload;
102103
auto payloadInSize = headerIn->payloadSize;
103-
if (!mCompressedData) { //we have raw data coming in from flp
104-
if (mVerbose) {
105-
LOG(info) << " parsing non compressed data in the data reader task with a payload of " << payloadInSize << " payload size";
106-
}
107-
mReader.setDataBuffer(payloadIn);
108-
mReader.setDataBufferSize(payloadInSize);
109-
mReader.configure(mByteSwap, mVerbose, mHeaderVerbose, mDataVerbose);
110-
if (mVerbose) {
111-
LOG(info) << "%%% about to run " << loopcounter << " %%%";
112-
}
113-
mReader.run();
114-
if (mVerbose) {
115-
LOG(info) << "%%% finished running " << loopcounter << " %%%";
116-
}
117-
loopcounter++;
118-
// mTracklets.insert(std::end(mTracklets), std::begin(mReader.getTracklets()), std::end(mReader.getTracklets()));
119-
// mCompressedDigits.insert(std::end(mCompressedDigits), std::begin(mReader.getCompressedDigits()), std::end(mReader.getCompressedDigits()));
120-
//mReader.clearall();
121-
if (mVerbose) {
122-
LOG(info) << "from parsing received: " << mTracklets.size() << " tracklets and " << mDigits.size() << " compressed digits";
123-
LOG(info) << "relevant vectors to read : " << mReader.sumTrackletsFound() << " tracklets and " << mReader.sumDigitsFound() << " compressed digits";
104+
if (std::string(headerIn->dataDescription.str) != std::string("DISTSUBTIMEFRAMEFLP")) {
105+
if (!mCompressedData) { //we have raw data coming in from flp
106+
if (mVerbose) {
107+
LOG(info) << " parsing non compressed data in the data reader task with a payload of " << payloadInSize << " payload size";
108+
}
109+
mReader.setDataBuffer(payloadIn);
110+
mReader.setDataBufferSize(payloadInSize);
111+
mReader.configure(mByteSwap, mVerbose, mHeaderVerbose, mDataVerbose);
112+
if (mVerbose) {
113+
LOG(info) << "%%% about to run " << loopcounter << " %%%";
114+
}
115+
mReader.run();
116+
if (mVerbose) {
117+
LOG(info) << "%%% finished running " << loopcounter << " %%%";
118+
}
119+
loopcounter++;
120+
// mTracklets.insert(std::end(mTracklets), std::begin(mReader.getTracklets()), std::end(mReader.getTracklets()));
121+
// mCompressedDigits.insert(std::end(mCompressedDigits), std::begin(mReader.getCompressedDigits()), std::end(mReader.getCompressedDigits()));
122+
//mReader.clearall();
123+
if (mVerbose) {
124+
LOG(info) << "from parsing received: " << mTracklets.size() << " tracklets and " << mDigits.size() << " compressed digits";
125+
LOG(info) << "relevant vectors to read : " << mReader.sumTrackletsFound() << " tracklets and " << mReader.sumDigitsFound() << " compressed digits";
126+
}
127+
// mTriggers = mReader.getIR();
128+
//get the payload of trigger and digits out.
129+
} else { // we have compressed data coming in from flp.
130+
mCompressedReader.setDataBuffer(payloadIn);
131+
mCompressedReader.setDataBufferSize(payloadInSize);
132+
mCompressedReader.configure(mByteSwap, mVerbose, mHeaderVerbose, mDataVerbose);
133+
mCompressedReader.run();
134+
//get the payload of trigger and digits out.
124135
}
125-
// mTriggers = mReader.getIR();
126-
//get the payload of trigger and digits out.
127-
} else { // we have compressed data coming in from flp.
128-
mCompressedReader.setDataBuffer(payloadIn);
129-
mCompressedReader.setDataBufferSize(payloadInSize);
130-
mCompressedReader.configure(mByteSwap, mVerbose, mHeaderVerbose, mDataVerbose);
131-
mCompressedReader.run();
132-
mTracklets = mCompressedReader.getTracklets();
133-
mDigits = mCompressedReader.getDigits();
134-
mTriggers = mCompressedReader.getIR();
135-
//get the payload of trigger and digits out.
136+
/* output */
137+
sendData(pc, false); //TODO do we ever have to not post the data. i.e. can we get here mid event? I dont think so.
138+
} else {
139+
sendData(pc, true);
136140
}
137141
}
138-
/* output */
139-
sendData(pc, false); //TODO do we ever have to not post the data. i.e. can we get here mid event? I dont think so.
140142
}
141143

142144
auto dataReadTime = std::chrono::high_resolution_clock::now() - dataReadStart;

0 commit comments

Comments
 (0)