Skip to content

Commit 37a39c9

Browse files
committed
o2 headers: adapt file source data to the new o2 DataHeader v2
1 parent 9fd5a40 commit 37a39c9

5 files changed

Lines changed: 158 additions & 63 deletions

File tree

src/common/SubTimeFrameFileReader.cxx

Lines changed: 140 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -69,51 +69,111 @@ void SubTimeFrameFileReader::visit(SubTimeFrame& pStf)
6969
}
7070
}
7171

72-
std::int64_t SubTimeFrameFileReader::getHeaderStackSize() // throws ios_base::failure
72+
std::size_t SubTimeFrameFileReader::getHeaderStackSize() // throws ios_base::failure
7373
{
74-
return sizeof(DataHeader); // TODO: allow arbitrary header stacks
74+
// Expect valid Stack in the file.
75+
// First Header must be DataHeader. The size is unknown since there are multiple versions.
76+
// Each header in the stack extends BaseHeader
7577

76-
#if 0
77-
std::int64_t lHdrStackSize = 0;
78+
// Read first the base header then the rest of the extended header. Keep going until the next flag is set.
79+
// reset the file pointer to the original incoming position, so the complete Stack can be read in
7880

79-
const auto lStartPos = mFile.tellg();
81+
bool readNextHeader = true;
82+
std::size_t lStackSize = 0;
83+
DataHeader lBaseHdr; // Use DataHeader since the BaseHeader has no default contructor.
8084

81-
DataHeader lBaseHdr;
85+
const auto lFilePosStart = position();
8286

83-
do {
84-
buffered_read(&lBaseHdr, sizeof(BaseHeader));
85-
if (nullptr == BaseHeader::get(reinterpret_cast<o2::byte*>(&lBaseHdr))) {
86-
// error: expected a header here
87-
return -1;
87+
const auto cMaxHeaders = 8; /* make sure we don't loop forever */
88+
auto lNumHeaders = 0;
89+
while (readNextHeader && (++lNumHeaders <= cMaxHeaders)) {
90+
buffered_read(&lBaseHdr, sizeof(BaseHeader)); // read BaseHeader only!
91+
92+
mFile.ignore(lBaseHdr.size());
93+
94+
lStackSize += lBaseHdr.size();
95+
readNextHeader = (lBaseHdr.next() != nullptr);
96+
}
97+
// reset the file pointer
98+
mFile.seekg(lFilePosStart);
99+
100+
if (lNumHeaders >= cMaxHeaders) {
101+
DDLOGF(fair::Severity::ERROR, "FileRead: Reached max number of headers allowed: {}.", cMaxHeaders);
102+
return 0;
103+
}
104+
105+
return lStackSize;
106+
}
107+
108+
Stack SubTimeFrameFileReader::getHeaderStack(std::size_t *pOrigsize) // throws ios_base::failure
109+
{
110+
const auto lStackSize = getHeaderStackSize();
111+
if (lStackSize < sizeof(BaseHeader)) {
112+
// error in the stream
113+
return Stack{};
114+
}
115+
116+
if (pOrigsize) {
117+
*pOrigsize = lStackSize;
118+
}
119+
120+
// std::unique_ptr<o2::byte[]> lStackMem = std::make_unique<o2::byte[]>(lStackSize);
121+
auto lStackMem = std::make_unique<o2::byte[]>(lStackSize);
122+
123+
// This must handle different versions of DataHeader
124+
buffered_read(lStackMem.get(), lStackSize);
125+
126+
// check if DataHeader needs an upgrade by looking at the version number
127+
const BaseHeader *lBaseOfDH = BaseHeader::get(lStackMem.get());
128+
if (!lBaseOfDH) {
129+
return Stack{};
130+
}
131+
132+
if (lBaseOfDH->headerVersion < DataHeader::sVersion) {
133+
DataHeader lNewDh;
134+
135+
// Write over the new DataHeader. We need to update some of the BaseHeader values.
136+
std::memcpy(&lNewDh, lBaseOfDH->data(), lBaseOfDH->size());
137+
// make sure to bump the version in the BaseHeader.
138+
// TODO: Is there a better way?
139+
lNewDh.headerSize = sizeof(DataHeader);
140+
lNewDh.headerVersion = DataHeader::sVersion;
141+
142+
if (lBaseOfDH->headerVersion == 1) {
143+
mDHUpdateFirstOrbit = true;
88144
} else {
89-
lHdrStackSize += lBaseHdr.headerSize;
145+
DDLOGF(fair::Severity::ERROR, "DataHeader version {} read from file is not upgraded to the current version {}",
146+
lBaseOfDH->headerVersion, DataHeader::sVersion);
147+
DDLOGF(fair::Severity::ERROR, "Try newer version of DataDistribution or file a BUG");
90148
}
91149

92-
// skip the rest of the current header
93-
if (lBaseHdr.headerSize > sizeof(BaseHeader)) {
94-
mFile.ignore(lBaseHdr.headerSize - sizeof(BaseHeader));
150+
assert (sizeof (DataHeader) > lBaseOfDH->size() ); // current DataHeader must be larger
151+
152+
if (lBaseOfDH->size() == lStackSize) {
153+
return Stack(lNewDh);
95154
} else {
96-
// error: invalid header size value
97-
return -1;
98-
}
99-
} while (lBaseHdr.next() != nullptr);
155+
assert(lBaseOfDH->size() < lStackSize);
100156

101-
// we should not eof here,
102-
if (mFile.eof())
103-
return -1;
157+
return Stack(
158+
lNewDh,
159+
Stack(lStackMem.get() + lBaseOfDH->size())
160+
);
161+
}
162+
}
104163

105-
// rewind the file to the start of Header stack
106-
mFile.seekg(lStartPos);
164+
return Stack(lStackMem.get());
165+
}
107166

108-
return lHdrStackSize;
109-
#endif
167+
Stack SubTimeFrameFileReader::getHeaderStack(SubTimeFrameFileBuilder &pFileBuilder) // throws ios_base::failure
168+
{
169+
// copy the memory to the SHM allocaton
170+
return Stack(pFileBuilder.getHeaderMemRes().allocator(), getHeaderStack());
110171
}
111172

173+
std::uint64_t SubTimeFrameFileReader::sStfId = 0; // TODO: add id to files metadata
174+
112175
std::unique_ptr<SubTimeFrame> SubTimeFrameFileReader::read(SubTimeFrameFileBuilder &pFileBuilder)
113176
{
114-
// TODO: add id to files metadata
115-
static std::uint64_t sStfId = 0;
116-
117177
// make sure headers and chunk pointers don't linger
118178
mStfData.clear();
119179

@@ -137,55 +197,75 @@ std::unique_ptr<SubTimeFrame> SubTimeFrameFileReader::read(SubTimeFrameFileBuild
137197
// NOTE: StfID will be updated from the stf header
138198
std::unique_ptr<SubTimeFrame> lStf = std::make_unique<SubTimeFrame>(sStfId++);
139199

140-
DataHeader lStfMetaDataHdr;
200+
std::unique_ptr<Stack> lMetaHdrStack;
201+
std::size_t lMetaHdrStackSize = 0;
202+
const DataHeader *lStfMetaDataHdr = nullptr;
141203
SubTimeFrameFileMeta lStfFileMeta;
142204

143205
try {
144206
// Read DataHeader + SubTimeFrameFileMeta
145-
buffered_read(&lStfMetaDataHdr, sizeof(DataHeader));
146-
buffered_read(&lStfFileMeta, sizeof(SubTimeFrameFileMeta));
207+
lMetaHdrStack = std::make_unique<Stack>(getHeaderStack(&lMetaHdrStackSize));
208+
lStfMetaDataHdr = o2::header::DataHeader::Get(lMetaHdrStack->first());
209+
if (!lStfMetaDataHdr) {
210+
DDLOGF(fair::Severity::ERROR, "Failed to read the TF file header. The file might be corrupted.");
211+
mFile.close();
212+
return nullptr;
213+
}
147214

215+
buffered_read(&lStfFileMeta, sizeof(SubTimeFrameFileMeta));
148216
} catch (const std::ios_base::failure& eFailExc) {
149-
DDLOG(fair::Severity::ERROR) << "Reading from file failed. Error: " << eFailExc.what();
217+
DDLOGF(fair::Severity::ERROR, "Reading from file failed. Error: {}", eFailExc.what());
218+
mFile.close();
150219
return nullptr;
151220
}
152221

153222
// verify we're actually reading the correct data in
154-
if (!(SubTimeFrameFileMeta::getDataHeader().dataDescription == lStfMetaDataHdr.dataDescription)) {
155-
DDLOG(fair::Severity::WARNING) << "Reading bad data: SubTimeFrame META header";
223+
if (!(SubTimeFrameFileMeta::getDataHeader().dataDescription == lStfMetaDataHdr->dataDescription)) {
224+
DDLOGF(fair::Severity::WARNING, "Reading bad data: SubTimeFrame META header");
156225
mFile.close();
157226
return nullptr;
158227
}
159228

160229
// prepare to read the TF data
161230
const auto lStfSizeInFile = lStfFileMeta.mStfSizeInFile;
162231
if (lStfSizeInFile == (sizeof(DataHeader) + sizeof(SubTimeFrameFileMeta))) {
163-
DDLOG(fair::Severity::WARNING) << "Reading an empty TF from file. Only meta information present";
232+
DDLOGF(fair::Severity::WARNING, "Reading an empty TF from file. Only meta information present");
233+
mFile.close();
164234
return nullptr;
165235
}
166236

167237
// check there's enough data in the file
168238
if ((lTfStartPosition + lStfSizeInFile) > this->size()) {
169-
DDLOG(fair::Severity::WARNING) << "Not enough data in file for this TF. Required: " << lStfSizeInFile
170-
<< ", available: " << (this->size() - lTfStartPosition);
239+
DDLOGF(fair::Severity::WARNING, "Not enough data in file for this TF. Required: {}, available: {}",
240+
lStfSizeInFile, (this->size() - lTfStartPosition));
171241
mFile.close();
172242
return nullptr;
173243
}
174244

175245
// Index
176246
// TODO: skip the index for now, check in future all data is there
177-
DataHeader lStfIndexHdr;
247+
std::unique_ptr<Stack> lStfIndexHdrStack;
248+
std::size_t lStfIndexHdrStackSize = 0;
249+
const DataHeader *lStfIndexHdr = nullptr;
178250
try {
179251
// Read DataHeader + SubTimeFrameFileMeta
180-
buffered_read(&lStfIndexHdr, sizeof(DataHeader));
181-
mFile.seekg(lStfIndexHdr.payloadSize, std::ios_base::cur);
252+
lStfIndexHdrStack = std::make_unique<Stack>(getHeaderStack(&lStfIndexHdrStackSize));
253+
lStfIndexHdr = o2::header::DataHeader::Get(lStfIndexHdrStack->first());
254+
if (!lStfIndexHdr) {
255+
DDLOGF(fair::Severity::ERROR, "Failed to read the TF index structure. The file might be corrupted.");
256+
return nullptr;
257+
}
258+
259+
mFile.ignore(lStfIndexHdr->payloadSize);
182260
} catch (const std::ios_base::failure& eFailExc) {
183-
DDLOG(fair::Severity::ERROR) << "Reading from file failed. Error: " << eFailExc.what();
261+
DDLOG(fair::Severity::ERROR) << "Reading TF index from file failed. Error: " << eFailExc.what();
184262
return nullptr;
185263
}
186264

187-
const auto lStfDataSize = lStfSizeInFile - (sizeof(DataHeader) + sizeof(SubTimeFrameFileMeta))
188-
- (sizeof (lStfIndexHdr) + lStfIndexHdr.payloadSize);
265+
// Remaining data size of the TF:
266+
// total size in file - meta (hdr+struct) - index (hdr + payload)
267+
const auto lStfDataSize = lStfSizeInFile - (lMetaHdrStackSize + sizeof(SubTimeFrameFileMeta))
268+
- (lStfIndexHdrStackSize + lStfIndexHdr->payloadSize);
189269

190270
// read all data blocks and headers
191271
assert(mStfData.empty());
@@ -195,31 +275,31 @@ std::unique_ptr<SubTimeFrame> SubTimeFrameFileReader::read(SubTimeFrameFileBuild
195275
// read <hdrStack + data> pairs
196276
while (lLeftToRead > 0) {
197277

198-
// read the header stack
199-
const std::int64_t lHdrSize = getHeaderStackSize();
200-
if (lHdrSize < std::int64_t(sizeof(DataHeader))) {
201-
// error while checking headers
202-
DDLOG(fair::Severity::WARNING) << "Reading bad data: Header stack cannot be parsed";
203-
mFile.close();
278+
// allocate and read the Headers
279+
std::size_t lDataHeaderStackSize = 0;
280+
Stack lDataHeaderStack = getHeaderStack(&lDataHeaderStackSize);
281+
const DataHeader *lDataHeader = o2::header::DataHeader::Get(lDataHeaderStack.first());
282+
if (!lDataHeader) {
283+
DDLOGF(fair::Severity::ERROR, "Failed to read the TF HBF DataHeader structure. The file might be corrupted.");
204284
return nullptr;
205285
}
206-
// allocate and read the Headers
207-
DataHeader lDataHeader;
208-
buffered_read(&lDataHeader, sizeof(DataHeader));
209286

210-
auto lHdrStackMsg = pFileBuilder.getHeaderMessage(lDataHeader, lStf->id());
287+
// TODO: fix fake first orbit
288+
reinterpret_cast<DataHeader*>(lDataHeaderStack.data())->firstTForbit = sStfId * 256;
289+
290+
auto lHdrStackMsg = pFileBuilder.getHeaderMessage(lDataHeaderStack, lStf->id());
211291
if (!lHdrStackMsg) {
212-
DDLOG(fair::Severity::WARNING) << "Out of memory: header message, allocation size: " << lHdrSize;
292+
DDLOGF(fair::Severity::WARNING, "Out of memory: header message, allocation size: {}", lDataHeaderStackSize);
213293
mFile.close();
214294
return nullptr;
215295
}
216296

217297
// read the data
218-
const std::uint64_t lDataSize = lDataHeader.payloadSize;
298+
const std::uint64_t lDataSize = lDataHeader->payloadSize;
219299

220300
auto lDataMsg = pFileBuilder.getDataMessage(lDataSize);
221301
if (!lDataMsg) {
222-
DDLOG(fair::Severity::WARNING) << "Out of memory: data message, allocation size: " << lDataSize;
302+
DDLOGF(fair::Severity::WARNING, "Out of memory: data message, allocation size: {}", lDataSize);
223303
mFile.close();
224304
return nullptr;
225305
}
@@ -228,10 +308,11 @@ std::unique_ptr<SubTimeFrame> SubTimeFrameFileReader::read(SubTimeFrameFileBuild
228308
mStfData.emplace_back(
229309
SubTimeFrame::StfData{
230310
std::move(lHdrStackMsg),
231-
std::move(lDataMsg) });
311+
std::move(lDataMsg) }
312+
);
232313

233314
// update the counter
234-
lLeftToRead -= (lHdrSize + lDataSize);
315+
lLeftToRead -= (lDataHeaderStackSize + lDataSize);
235316
}
236317

237318
if (lLeftToRead < 0) {

src/common/include/SubTimeFrameBuilder.h

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -78,24 +78,26 @@ class SubTimeFrameFileBuilder
7878
}
7979

8080
// allocate appropriate message for the header
81-
FairMQMessagePtr getHeaderMessage(const o2::header::DataHeader &pDh, const std::uint64_t pTfId) {
81+
FairMQMessagePtr getHeaderMessage(const o2::header::Stack &pIncomingStack, const std::uint64_t pTfId) {
8282
std::unique_ptr<FairMQMessage> lMsg;
8383

8484
if (mDplEnabled) {
8585
auto lStack = o2::header::Stack(mHeaderMemRes->allocator(),
86-
pDh,
86+
pIncomingStack,
8787
o2::framework::DataProcessingHeader{pTfId}
8888
);
8989

9090
lMsg = mHeaderMemRes->NewFairMQMessageFromPtr(lStack.data());
9191
} else {
92-
auto lHdrMsgStack = o2::header::Stack(mHeaderMemRes->allocator(), pDh);
92+
auto lHdrMsgStack = o2::header::Stack(mHeaderMemRes->allocator(), pIncomingStack);
9393
lMsg = mHeaderMemRes->NewFairMQMessageFromPtr(lHdrMsgStack.data());
9494
}
9595

9696
return lMsg;
9797
}
9898

99+
auto& getHeaderMemRes() const { return *mHeaderMemRes; }
100+
99101
private:
100102

101103
bool mDplEnabled;

src/common/include/SubTimeFrameDPL.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,8 @@ class StfToDplAdapter : public ISubTimeFrameVisitor
5656

5757
class DplToStfAdapter : public ISubTimeFrameVisitor
5858
{
59+
constexpr static std::uint64_t sCurrentTfId = 0;
60+
5961
public:
6062
DplToStfAdapter() = default;
6163
virtual ~DplToStfAdapter() = default;

src/common/include/SubTimeFrameDataModel.h

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -232,6 +232,7 @@ class SubTimeFrame : public IDataModelObject
232232

233233
inline const o2hdr::DataHeader& getDataHeader() const
234234
{
235+
// this is fine since we created the DataHeader there
235236
return *reinterpret_cast<o2hdr::DataHeader*>(mHeader->GetData());
236237
}
237238

src/common/include/SubTimeFrameFileReader.h

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
#include "SubTimeFrameDataModel.h"
1818
#include <Headers/DataHeader.h>
19+
#include <Headers/Stack.h>
1920

2021
#include <boost/filesystem.hpp>
2122
#include <fstream>
@@ -68,10 +69,18 @@ class SubTimeFrameFileReader : public ISubTimeFrameVisitor
6869
return mFile.read(reinterpret_cast<char*>(pPtr), pLen);
6970
}
7071

71-
std::int64_t getHeaderStackSize();
72+
std::size_t getHeaderStackSize();
73+
o2::header::Stack getHeaderStack(std::size_t *pOrigsize = nullptr);
74+
o2::header::Stack getHeaderStack(SubTimeFrameFileBuilder &pFileBuilder);
7275

7376
// vector of <hdr, fmqMsg> elements of a tf read from the file
7477
std::vector<SubTimeFrame::StfData> mStfData;
78+
79+
80+
// flags for upgrading DataHeader versions
81+
bool mDHUpdateFirstOrbit = false; // Set first orbit of in the new DataHeader
82+
static std::uint64_t sStfId; // TODO: add id to files metadata
83+
7584
};
7685
}
7786
} /* o2::DataDistribution */

0 commit comments

Comments
 (0)