diff --git a/artdaq/ArtModules/ArtdaqRunInfoServiceInterface.h b/artdaq/ArtModules/ArtdaqRunInfoServiceInterface.h new file mode 100644 index 00000000..14a5e96c --- /dev/null +++ b/artdaq/ArtModules/ArtdaqRunInfoServiceInterface.h @@ -0,0 +1,83 @@ +#ifndef artdaq_ArtModules_ArtdaqRunInfoServiceInterface_h +#define artdaq_ArtModules_ArtdaqRunInfoServiceInterface_h + +#include "art/Framework/Services/Registry/ServiceDeclarationMacros.h" +#include "canvas/Persistency/Provenance/IDNumber.h" + +#include +#include +#include +#include + +#include +#include +#include + +class ArtdaqRunInfoServiceInterface +{ +public: + ArtdaqRunInfoServiceInterface() = default; + + virtual ~ArtdaqRunInfoServiceInterface() = default; + + virtual bool addSubrunRecord( + art::RunNumber_t run, + art::SubRunNumber_t subrun, + size_t nEvents, + art::EventNumber_t firstEvent, + art::EventNumber_t lastEvent, + std::string const& datastream = "default") = 0; + + virtual bool addFileSummary( + std::string const& fileName, + art::RunNumber_t run, + art::SubRunNumber_t firstSubrun, + art::SubRunNumber_t lastSubrun, + size_t nEvents, + size_t fileSize, + std::string const& metadata = "{}", + std::string const& datastream = "default") = 0; + +protected: + bool appendToCsv(std::string const& path, std::string const& header, std::string const& row) + { + int fd = open(path.c_str(), O_WRONLY | O_CREAT | O_APPEND, 0666); + if (fd < 0) { return false; } + + flock(fd, LOCK_EX); + + std::string buf; + if (lseek(fd, 0, SEEK_END) == 0) { buf = header; } + buf += row; + + ssize_t total = 0; + auto remaining = static_cast(buf.size()); + while (remaining > 0) + { + ssize_t written = ::write(fd, buf.c_str() + total, static_cast(remaining)); + if (written < 0) + { + if (errno == EINTR) { continue; } + flock(fd, LOCK_UN); + close(fd); + return false; + } + total += written; + remaining -= written; + } + + flock(fd, LOCK_UN); + close(fd); + return true; + } + +private: + ArtdaqRunInfoServiceInterface(ArtdaqRunInfoServiceInterface const&) = delete; + ArtdaqRunInfoServiceInterface(ArtdaqRunInfoServiceInterface&&) = delete; + ArtdaqRunInfoServiceInterface& operator=(ArtdaqRunInfoServiceInterface const&) = delete; + ArtdaqRunInfoServiceInterface& operator=(ArtdaqRunInfoServiceInterface&&) = delete; +}; + +DECLARE_ART_SERVICE_INTERFACE(ArtdaqRunInfoServiceInterface, LEGACY) + +#endif /* artdaq_ArtModules_ArtdaqRunInfoServiceInterface_h */ diff --git a/artdaq/ArtModules/ArtdaqRunInfoService_service.cc b/artdaq/ArtModules/ArtdaqRunInfoService_service.cc new file mode 100644 index 00000000..99c3c19f --- /dev/null +++ b/artdaq/ArtModules/ArtdaqRunInfoService_service.cc @@ -0,0 +1,103 @@ +#include "TRACE/tracemf.h" + +#include "artdaq/ArtModules/ArtdaqRunInfoServiceInterface.h" + +#include "art/Framework/Services/Registry/ServiceDefinitionMacros.h" + +#include "fhiclcpp/ParameterSet.h" + +#include +#include +#include +#include +#include + +#define TRACE_NAME "ArtdaqRunInfoService" + +class ArtdaqRunInfoService : public ArtdaqRunInfoServiceInterface +{ +public: + ArtdaqRunInfoService(fhicl::ParameterSet const& pset, art::ActivityRegistry&); + ~ArtdaqRunInfoService() override = default; + + bool addSubrunRecord( + art::RunNumber_t run, + art::SubRunNumber_t subrun, + size_t nEvents, + art::EventNumber_t firstEvent, + art::EventNumber_t lastEvent, + std::string const& datastream) override; + + bool addFileSummary( + std::string const& fileName, + art::RunNumber_t run, + art::SubRunNumber_t firstSubrun, + art::SubRunNumber_t lastSubrun, + size_t nEvents, + size_t fileSize, + std::string const& metadata, + std::string const& datastream) override; + +private: + std::string summaryDir_; +}; + +DECLARE_ART_SERVICE_INTERFACE_IMPL(ArtdaqRunInfoService, ArtdaqRunInfoServiceInterface, LEGACY) + +ArtdaqRunInfoService::ArtdaqRunInfoService(fhicl::ParameterSet const& pset, art::ActivityRegistry& /*unused*/) + : summaryDir_(pset.get("summaryDir", "")) +{ + TLOG(TLVL_INFO) << "ArtdaqRunInfoService: summaryDir=\"" << summaryDir_ << "\""; +} + +bool ArtdaqRunInfoService::addSubrunRecord( + art::RunNumber_t run, + art::SubRunNumber_t subrun, + size_t nEvents, + art::EventNumber_t firstEvent, + art::EventNumber_t lastEvent, + std::string const& datastream) +{ + if (summaryDir_.empty()) { return true; } + + std::ostringstream fname; + fname << summaryDir_; + if (summaryDir_.back() != '/') { fname << '/'; } + fname << "subrun_record_run" << std::setw(6) << std::setfill('0') << run << ".csv"; + + std::ostringstream row; + row << run << "," << subrun << "," << nEvents << "," << firstEvent << "," << lastEvent + << "," << datastream << "\n"; + + return appendToCsv(fname.str(), + "run,subrun,n_events,first_event,last_event,datastream\n", row.str()); +} + +bool ArtdaqRunInfoService::addFileSummary( + std::string const& fileName, + art::RunNumber_t run, + art::SubRunNumber_t firstSubrun, + art::SubRunNumber_t lastSubrun, + size_t nEvents, + size_t fileSize, + std::string const& /*metadata*/, + std::string const& datastream) +{ + if (summaryDir_.empty()) { return true; } + + std::ostringstream fname; + fname << summaryDir_; + if (summaryDir_.back() != '/') { fname << '/'; } + fname << "file_summary_run" << std::setw(6) << std::setfill('0') << run << ".csv"; + + std::string const outputFile = std::filesystem::path(fileName).filename().string(); + + std::ostringstream row; + row << outputFile << "," << run << "," << firstSubrun << "," << lastSubrun + << "," << nEvents << "," << fileSize << "," << datastream << "\n"; + + return appendToCsv(fname.str(), + "file_name,run,first_subrun,last_subrun,n_events,file_size,datastream\n", row.str()); +} + +DEFINE_ART_SERVICE_INTERFACE_IMPL(ArtdaqRunInfoService, ArtdaqRunInfoServiceInterface) diff --git a/artdaq/ArtModules/CMakeLists.txt b/artdaq/ArtModules/CMakeLists.txt index f6894d42..de11edb0 100644 --- a/artdaq/ArtModules/CMakeLists.txt +++ b/artdaq/ArtModules/CMakeLists.txt @@ -54,6 +54,7 @@ cet_build_plugin(RootDAQOutMF art::module LIBRARIES REG artdaq::RootDAQOutput artdaq::ArtModules + artdaq_plugin_types::ArtdaqRunInfoService TRACE::MF art_root_io::detail art_root_io::art_root_io @@ -128,6 +129,21 @@ cet_make_library(SOURCE art::Framework_Services_Registry ) +cet_make_library(LIBRARY_NAME ArtdaqRunInfoService INTERFACE + EXPORT_SET AMPluginTypes + SOURCE ArtdaqRunInfoServiceInterface.h + LIBRARIES INTERFACE + canvas::canvas + art_plugin_types::serviceDeclaration +) + +cet_build_plugin(ArtdaqRunInfoService art::service + LIBRARIES PRIVATE + artdaq_plugin_types::ArtdaqRunInfoService + fhiclcpp::fhiclcpp + TRACE::MF +) + cet_make_library(LIBRARY_NAME ArtdaqFragmentNamingService INTERFACE EXPORT_SET AMPluginTypes SOURCE ArtdaqFragmentNamingService.h diff --git a/artdaq/ArtModules/RootDAQOutMF_module.cc b/artdaq/ArtModules/RootDAQOutMF_module.cc index 76094dc7..d7405a31 100644 --- a/artdaq/ArtModules/RootDAQOutMF_module.cc +++ b/artdaq/ArtModules/RootDAQOutMF_module.cc @@ -4,6 +4,7 @@ #include "artdaq/DAQdata/Globals.hh" #define TRACE_NAME (app_name + "_RootDAQOutMF").c_str() +#include "artdaq/ArtModules/ArtdaqRunInfoServiceInterface.h" #include "artdaq/ArtModules/ArtdaqSharedMemoryServiceInterface.h" #include "artdaq/ArtModules/RootDAQOutFile.h" @@ -41,12 +42,9 @@ #include "fhiclcpp/types/TableFragment.h" #include "messagefacility/MessageLogger/MessageLogger.h" -#include -#include #include + #include -#include -#include #include #include #include @@ -115,87 +113,6 @@ struct SubrunStats art::EventNumber_t lastEvent{0}; }; -static void writeSummaryFile( - std::string const& summaryDir, - std::map const& subrunStats, - std::string const& closedFileName) -{ - if (summaryDir.empty() || subrunStats.empty()) { return; } - - std::string const outputFile = std::filesystem::path(closedFileName).filename().string(); - - // Group rows by run number so that each run gets its own CSV file, - // even when an output file spans multiple runs. - std::map runContents; - for (auto const& [srid, stats] : subrunStats) - { - runContents[srid.run()] - << srid.run() - << "," << srid.subRun() - << "," << stats.nEvents - << "," << stats.firstEvent - << "," << stats.lastEvent - << "," << outputFile - << "\n"; - } - - // One CSV file per run, appended — matches CFODataReceiver convention - for (auto const& [run, contentStream] : runContents) - { - std::ostringstream fname; - fname << summaryDir; - if (summaryDir.back() != '/') { fname << '/'; } - fname << "subrun_record_run" << std::setw(6) << std::setfill('0') << run << ".csv"; - - int fd = open(fname.str().c_str(), O_WRONLY | O_CREAT | O_APPEND, 0666); - if (fd < 0) - { - TLOG(TLVL_WARNING) << "writeSummaryFile: could not open \"" << fname.str() << "\" for writing: " << strerror(errno); - continue; - } - - flock(fd, LOCK_EX); - - std::string const rows = contentStream.str(); - std::string buf; - if (lseek(fd, 0, SEEK_END) == 0) - { - buf = "run,subrun,n_events,first_event,last_event,output_file\n"; - } - buf += rows; - - // Loop to handle partial writes and EINTR - if (buf.size() > static_cast(std::numeric_limits::max())) - { - TLOG(TLVL_ERROR) << "writeSummaryFile: buffer too large (" << buf.size() << " bytes) to write to \"" - << fname.str() << "\""; - } - else - { - ssize_t total = 0; - auto remaining = static_cast(buf.size()); - while (remaining > 0) - { - ssize_t written = ::write(fd, buf.c_str() + total, static_cast(remaining)); - if (written < 0) - { - if (errno == EINTR) { continue; } - TLOG(TLVL_ERROR) << "writeSummaryFile: write error to \"" << fname.str() << "\": " << strerror(errno); - break; - } - total += written; - remaining -= written; - } - } - - flock(fd, LOCK_UN); - close(fd); - - size_t const rowsWritten = static_cast(std::count(rows.begin(), rows.end(), '\n')); - TLOG(TLVL_DEBUG) << "writeSummaryFile: appended " << rowsWritten - << " subrun row(s) to \"" << fname.str() << "\""; - } -} } // namespace namespace art { @@ -265,7 +182,17 @@ class RootDAQOutMF final : public OutputModule fhicl::Sequence> replacementList{fhicl::Name("replacementList")}; }; fhicl::OptionalSequence> fileNameSubstitutions{Name("fileNameSubstitutions")}; - Atom summaryDir{Name("subrunRecordDir"), Comment("Directory for per-file CSV subrun record (subrun/event statistics). Empty = disabled."), ""}; + Atom datastream{Name("datastream"), + Comment("Datastream label included in summary records written by\n" + "ArtdaqRunInfoServiceInterface (e.g. \"triggered\", \"lumistream\").\n" + "The default implementation (ArtdaqRunInfoService) writes CSV files;\n" + "experiments can provide their own (e.g. Postgres).\n" + "Enable the service in FHiCL:\n" + " services.ArtdaqRunInfoServiceInterface: {\n" + " service_provider: \"ArtdaqRunInfoService\"\n" + " summaryDir: \"/path/to/output\"\n" + " }"), + "default"}; Config() { @@ -431,7 +358,8 @@ class RootDAQOutMF final : public OutputModule // ParameterSet information in the downstream file, such as when mixing. bool writeParameterSets_; ClosingCriteria fileProperties_; - string summaryDir_; + string datastream_; + ArtdaqRunInfoServiceInterface* runInfoService_{nullptr}; size_t filesOpenedInRun_{0}; size_t filesClosedInRun_{0}; // Shared %# sequence counter across all OutputFileBundle instances. @@ -505,7 +433,7 @@ RootDAQOutMF::RootDAQOutMF(Parameters const& config) , dropMetaDataForDroppedData_{config().dropMetaDataForDroppedData()} , writeParameterSets_{config().writeParameterSets()} , fileProperties_{config().fileProperties()} - , summaryDir_{config().summaryDir()} + , datastream_{config().datastream()} , rpm_{config.get_PSet()} { TLOG(TLVL_INFO) << "RootDAQOutMF_module (s124 version) CONSTRUCTOR Start this=" << static_cast(this) @@ -555,6 +483,20 @@ RootDAQOutMF::RootDAQOutMF(Parameters const& config) "problems\n" << "with analysis reproducibility.\n"; } + + // Probe for ArtdaqRunInfoServiceInterface — writes per-subrun and per-file + // summary records at file close. See the "datastream" Config comment above. + try + { + art::ServiceHandle svc; + runInfoService_ = &*svc; + TLOG(TLVL_INFO) << "RootDAQOutMF: ArtdaqRunInfoServiceInterface available"; + } + catch (art::Exception const&) + { + runInfoService_ = nullptr; + TLOG(TLVL_INFO) << "RootDAQOutMF: ArtdaqRunInfoServiceInterface not configured, summary records disabled"; + } } void RootDAQOutMF::openFile(FileBlock const& fb) @@ -962,7 +904,39 @@ void RootDAQOutMF::closePendingFile(std::unique_ptr& bundle) bundle->file.reset(); bundle->closedFileName = fileNameAtClose(bundle->fRenamer, bundle->tmpFileName); - writeSummaryFile(summaryDir_, bundle->subrunStats, bundle->closedFileName); + if (runInfoService_ && !bundle->subrunStats.empty()) + { + for (auto const& [srid, stats] : bundle->subrunStats) + { + runInfoService_->addSubrunRecord( + srid.run(), srid.subRun(), + stats.nEvents, stats.firstEvent, stats.lastEvent, + datastream_); + } + + art::SubRunNumber_t firstSubrun = bundle->subrunStats.begin()->first.subRun(); + art::SubRunNumber_t lastSubrun = bundle->subrunStats.rbegin()->first.subRun(); + size_t totalEvents = 0; + for (auto const& [_, stats] : bundle->subrunStats) + totalEvents += stats.nEvents; + + size_t fileSize = 0; + std::error_code ec; + auto sz = std::filesystem::file_size(bundle->closedFileName, ec); + if (!ec) fileSize = static_cast(sz); + + art::RunNumber_t run = bundle->subrunStats.begin()->first.run(); + + std::string dirPath = std::filesystem::path(bundle->closedFileName).parent_path().string(); + char hostname[256] = {}; + gethostname(hostname, sizeof(hostname) - 1); + std::string metadata = "{\"path\":\"" + dirPath + "\",\"hostname\":\"" + hostname + "\"}"; + + runInfoService_->addFileSummary( + bundle->closedFileName, run, + firstSubrun, lastSubrun, totalEvents, fileSize, + metadata, datastream_); + } ++filesClosedInRun_; TLOG(TLVL_DEBUG) << __func__ << ": filesClosedInRun_ now " << filesClosedInRun_ << ", metricMan=" << (metricMan ? "non-null" : "NULL");