Skip to content

Commit 4f77b72

Browse files
committed
Track timeframe reading progress in reader state machine
1 parent 4601255 commit 4f77b72

1 file changed

Lines changed: 63 additions & 53 deletions

File tree

Framework/AnalysisSupport/src/AODJAlienReaderHelpers.cxx

Lines changed: 63 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -253,80 +253,88 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const
253253
*numTF = ntf;
254254
};
255255
enum class TFReaderState {
256-
READ_TIMEFRAME,
256+
READ_FIRST_TABLE,
257+
READ_FIRST_TABLE_FROM_NEXT_FILE,
258+
READ_NEXT_TABLE,
257259
TRY_NEXT_FILE,
258260
TIMEFRAME_READ,
259261
INVALID_TIMEFRAME,
260262
END_OF_INPUT,
261263
};
262-
auto readState = TFReaderState::READ_TIMEFRAME;
263-
bool triedNextFile = false;
264-
while (readState == TFReaderState::READ_TIMEFRAME || readState == TFReaderState::TRY_NEXT_FILE) {
264+
auto readState = TFReaderState::READ_FIRST_TABLE;
265+
size_t routeIndex = 0;
266+
while (readState == TFReaderState::READ_FIRST_TABLE ||
267+
readState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE ||
268+
readState == TFReaderState::READ_NEXT_TABLE ||
269+
readState == TFReaderState::TRY_NEXT_FILE) {
265270
if (readState == TFReaderState::TRY_NEXT_FILE) {
266271
fcnt += device.maxInputTimeslices;
267272
if (didir->atEnd(fcnt)) {
268273
readState = TFReaderState::END_OF_INPUT;
269274
break;
270275
}
271276
ntf = 0;
272-
triedNextFile = true;
277+
routeIndex = 0;
278+
readState = TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE;
273279
}
274280

275-
readState = TFReaderState::TIMEFRAME_READ;
276-
bool firstTable = true;
277-
for (auto& route : requestedTables) {
278-
if ((device.inputTimesliceId % route.maxTimeslices) != route.timeslice) {
279-
continue;
280-
}
281+
while (routeIndex < requestedTables.size() &&
282+
(device.inputTimesliceId % requestedTables[routeIndex].maxTimeslices) != requestedTables[routeIndex].timeslice) {
283+
++routeIndex;
284+
}
285+
if (routeIndex == requestedTables.size()) {
286+
readState = TFReaderState::TIMEFRAME_READ;
287+
break;
288+
}
281289

282-
// create header
283-
auto concrete = DataSpecUtils::asConcreteDataMatcher(route.matcher);
284-
auto dh = header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
285-
bool wasAOD = std::ranges::any_of(route.matcher.metadata, [](ConfigParamSpec const& p) { return p.name.starts_with("aod-origin-replaced"); });
286-
287-
try {
288-
if (!didir->readTree(outputs, dh, fcnt, ntf, totalSizeCompressed, totalSizeUncompressed, wasAOD)) {
289-
if (firstTable && !triedNextFile) {
290-
readState = TFReaderState::TRY_NEXT_FILE;
291-
break;
292-
}
293-
LOGP(fatal, "Can not retrieve tree for table {}: fileCounter {}, timeFrame {}", concrete.origin.as<std::string>(), fcnt, ntf);
294-
throw std::runtime_error("Processing is stopped!");
295-
}
296-
} catch (InvalidAODReadError const& e) {
297-
if (!skipInvalidReads) {
298-
throw;
290+
auto& route = requestedTables[routeIndex];
291+
auto concrete = DataSpecUtils::asConcreteDataMatcher(route.matcher);
292+
auto dh = header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
293+
bool wasAOD = std::ranges::any_of(route.matcher.metadata, [](ConfigParamSpec const& p) { return p.name.starts_with("aod-origin-replaced"); });
294+
295+
try {
296+
if (!didir->readTree(outputs, dh, fcnt, ntf, totalSizeCompressed, totalSizeUncompressed, wasAOD)) {
297+
if (readState == TFReaderState::READ_FIRST_TABLE) {
298+
readState = TFReaderState::TRY_NEXT_FILE;
299+
continue;
299300
}
300-
skipInvalidRead(concrete.origin, e);
301-
readState = TFReaderState::INVALID_TIMEFRAME;
302-
break;
301+
LOGP(fatal, "Can not retrieve tree for table {}: fileCounter {}, timeFrame {}", concrete.origin.as<std::string>(), fcnt, ntf);
302+
throw std::runtime_error("Processing is stopped!");
303+
}
304+
} catch (InvalidAODReadError const& e) {
305+
if (!skipInvalidReads) {
306+
throw;
303307
}
308+
skipInvalidRead(concrete.origin, e);
309+
readState = TFReaderState::INVALID_TIMEFRAME;
310+
break;
311+
}
304312

305-
if (firstTable) {
306-
if (reportTFN) {
307-
// TF number
308-
auto timeFrameNumber = didir->getTimeFrameNumber(dh, fcnt, ntf);
309-
auto o = Output(TFNumberHeader);
310-
outputs.make<uint64_t>(o) = timeFrameNumber;
311-
}
313+
if (readState == TFReaderState::READ_FIRST_TABLE || readState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE) {
314+
if (reportTFN) {
315+
// TF number
316+
auto timeFrameNumber = didir->getTimeFrameNumber(dh, fcnt, ntf);
317+
auto o = Output(TFNumberHeader);
318+
outputs.make<uint64_t>(o) = timeFrameNumber;
319+
}
312320

313-
if (reportTFFileName) {
314-
// Origin file name for derived output map
315-
auto o2 = Output(TFFileNameHeader);
316-
auto fileAndFolder = didir->getFileFolder(dh, fcnt, ntf);
317-
auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(fileAndFolder.filesystem());
318-
auto* f = dynamic_cast<TFile*>(rootFS->GetFile());
319-
std::string currentFilename(f->GetFile()->GetName());
320-
if (strcmp(f->GetEndpointUrl()->GetProtocol(), "file") == 0 && f->GetEndpointUrl()->GetFile()[0] != '/') {
321-
// This is not an absolute local path. Make it absolute.
322-
static std::string pwd = gSystem->pwd() + std::string("/");
323-
currentFilename = pwd + std::string(f->GetName());
324-
}
325-
outputs.make<std::string>(o2) = currentFilename;
321+
if (reportTFFileName) {
322+
// Origin file name for derived output map
323+
auto o2 = Output(TFFileNameHeader);
324+
auto fileAndFolder = didir->getFileFolder(dh, fcnt, ntf);
325+
auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(fileAndFolder.filesystem());
326+
auto* f = dynamic_cast<TFile*>(rootFS->GetFile());
327+
std::string currentFilename(f->GetFile()->GetName());
328+
if (strcmp(f->GetEndpointUrl()->GetProtocol(), "file") == 0 && f->GetEndpointUrl()->GetFile()[0] != '/') {
329+
// This is not an absolute local path. Make it absolute.
330+
static std::string pwd = gSystem->pwd() + std::string("/");
331+
currentFilename = pwd + std::string(f->GetName());
326332
}
333+
outputs.make<std::string>(o2) = currentFilename;
327334
}
328-
firstTable = false;
329335
}
336+
++routeIndex;
337+
readState = TFReaderState::READ_NEXT_TABLE;
330338
}
331339

332340
switch (readState) {
@@ -341,7 +349,9 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const
341349
control.endOfStream();
342350
control.readyToQuit(QuitRequest::Me);
343351
return;
344-
case TFReaderState::READ_TIMEFRAME:
352+
case TFReaderState::READ_FIRST_TABLE:
353+
case TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE:
354+
case TFReaderState::READ_NEXT_TABLE:
345355
case TFReaderState::TRY_NEXT_FILE:
346356
throw std::logic_error("Invalid timeframe read state");
347357
}

0 commit comments

Comments
 (0)