8#include "G4Threading.hh"
10std::mutex GRunAction::completed_run_data_mutex;
11GRunAction::CompletedRunData GRunAction::completed_worker_run_data;
17 std::shared_ptr<gdynamicdigitization::dRoutinesMap> digi_map,
18 std::shared_ptr<GAnalysisAccumulator> analyzer,
21 digitization_routines_map(std::move(digi_map)),
22 analysis_accumulator(std::move(analyzer)),
23 sro_factory(std::move(sro)) {
24 const auto desc = std::to_string(G4Threading::G4GetThreadId());
30G4Run *GRunAction::GenerateRun() {
38void GRunAction::BeginOfRunAction(
const G4Run *aRun) {
39 const auto thread_id = G4Threading::G4GetThreadId();
40 const auto run = aRun->GetRunID();
42 auto run_header = std::make_unique<GRunHeader>(goptions, run, thread_id);
43 run_data = std::make_unique<GRunDataCollection>(goptions, std::move(run_header));
44 if (analysis_accumulator !=
nullptr) {
45 analysis_run_number = analysis_accumulator->currentRunNumber();
46 analysis_shard = std::make_unique<GAnalysisShard>();
49 const auto neventsThisRun = aRun->GetNumberOfEventToBeProcessed();
50 if (sro_factory && IsMaster()) {
51 try { sro_factory->begin_run(run,
static_cast<GSROEventId>(neventsThisRun)); }
52 catch (
const std::exception& error) {
58 need_a_thread_streamer =
false;
59 need_a_run_streamer =
false;
63 if (digitization_routines_map !=
nullptr) {
64 for (
const auto &[plugin, digiRoutine]: *digitization_routines_map) {
65 if (digiRoutine ==
nullptr) {
67 " null digitization routine registered for plugin ", plugin);
70 if (digiRoutine->collection_mode() == CollectionMode::event) {
71 need_a_thread_streamer =
true;
72 }
else if (digiRoutine->collection_mode() == CollectionMode::run) {
73 to_normalize[plugin] = digiRoutine->variables_to_normalize();
74 need_a_run_streamer =
true;
79 " digitization_routines_map is null - streamer mode detection skipped.");
83 if (streamer_definition.type ==
"event") {
84 need_a_thread_streamer =
true;
92 if (!IsMaster() && need_a_thread_streamer) {
93 if (gstreamer_threads_map ==
nullptr) {
94 log->info(1,
"Defining thread gstreamers for run ", run,
" in thread ", thread_id);
98 if (gstreamer_threads_map ==
nullptr) {
99 log->error(1, FUNCTION_NAME,
" gstreamer_threads_map is null in thread ", thread_id,
100 " - cannot open connections.");
103 for (
const auto &[name, gstreamer]: *gstreamer_threads_map) {
104 if (gstreamer ==
nullptr) {
106 "Null GStreamer entry ", name,
" in thread ", thread_id);
109 if (!gstreamer->openConnection()) {
111 "Failed to open connection for GStreamer ", name,
112 " in thread ", thread_id);
115 log->info(2, FUNCTION_NAME,
"Worker thread [", thread_id,
"]: opening connection for ",
117 " for run ", run,
". Number of events to be processed: ", neventsThisRun);
122 else if (IsMaster() && need_a_run_streamer) {
123 if (gstreamer_run_map ==
nullptr) {
124 log->info(1,
"Defining run gstreamers for run ", run);
128 if (gstreamer_run_map ==
nullptr) {
129 log->error(1, FUNCTION_NAME,
" gstreamer_run_map is null in master thread ",
130 " - cannot open connections.");
133 for (
const auto &[name, gstreamer]: *gstreamer_run_map) {
134 if (gstreamer ==
nullptr) {
136 "Null GStreamer entry ", name,
" in master thread");
139 if (!gstreamer->openConnection()) {
141 "Failed to open connection for GStreamer in master thread ", name);
144 log->info(2, FUNCTION_NAME,
"Master Thread: opening connection for ",
146 " for run ", run,
". Number of events to be processed: ", neventsThisRun);
153void GRunAction::EndOfRunAction(
const G4Run *aRun) {
154 const auto thread_id = G4Threading::G4GetThreadId();
155 const auto runNumber = aRun->GetRunID();
156 const std::string what_am_i = IsMaster() ?
"Master" :
"Worker";
158 if (sro_factory && IsMaster()) {
159 try { sro_factory->finish_run(aRun->GetNumberOfEvent() == aRun->GetNumberOfEventToBeProcessed()); }
160 catch (
const std::exception& error) {
166 if (!IsMaster() && need_a_thread_streamer) {
167 if (gstreamer_threads_map ==
nullptr) {
169 " gstreamer_map is null in thread ", thread_id,
170 " - cannot close connections.");
172 for (
const auto &[name, gstreamer]: *gstreamer_threads_map) {
173 log->info(2, FUNCTION_NAME,
" ", what_am_i,
" [", thread_id,
"], for run ", runNumber,
174 " closing connection for gstreamer ", name);
176 if (gstreamer ==
nullptr) {
178 "Null GStreamer entry ", name,
" in thread ", thread_id);
181 if (!gstreamer->closeConnection()) {
182 log->error(1,
"Failed to close connection for GStreamer ", name,
" in thread ", thread_id);
189 if (analysis_accumulator !=
nullptr && analysis_shard !=
nullptr) {
190 if (!analysis_shard->empty()) { analysis_accumulator->merge(std::move(*analysis_shard)); }
191 analysis_shard.reset();
199 if (need_a_run_streamer) { stash_worker_run_data(); }
203 if (IsMaster() && need_a_run_streamer) {
206 auto completed_run_data = take_completed_worker_run_data();
207 log->info(2, FUNCTION_NAME,
208 " master collected ",
static_cast<int>(completed_run_data.size()),
209 " worker run_data object(s) for run ", runNumber);
211 std::shared_ptr<GRunDataCollection> merged_run_data;
213 for (
auto &worker_run_data: completed_run_data) {
214 if (worker_run_data ==
nullptr) {
220 if (merged_run_data ==
nullptr) {
221 auto merged_header = std::make_unique<GRunHeader>(goptions, runNumber, thread_id);
222 merged_run_data = std::make_shared<GRunDataCollection>(goptions, std::move(merged_header));
225 merged_run_data->merge(*worker_run_data);
229 if (merged_run_data !=
nullptr) {
230 publish_run_data(merged_run_data);
233 if (gstreamer_run_map ==
nullptr) {
235 " gstreamer_map is null in master thread - cannot close connections.");
238 for (
const auto &[name, gstreamer]: *gstreamer_run_map) {
239 log->info(2, FUNCTION_NAME,
" ", what_am_i,
" for run ", runNumber,
240 " closing connection for gstreamer ", name);
242 if (gstreamer ==
nullptr) {
244 "Null GStreamer entry ", name,
" in master thread");
247 if (!gstreamer->closeConnection()) {
248 log->error(1,
"Failed to close connection for GStreamer ", name,
" in master thread");
257void GRunAction::stash_worker_run_data() {
258 if (run_data ==
nullptr) {
262 std::scoped_lock lock(completed_run_data_mutex);
263 completed_worker_run_data.emplace_back(std::move(run_data));
267auto GRunAction::take_completed_worker_run_data() -> CompletedRunData {
268 std::scoped_lock lock(completed_run_data_mutex);
270 auto result = std::move(completed_worker_run_data);
271 completed_worker_run_data.clear();
278void GRunAction::publish_run_data(
const std::shared_ptr<GRunDataCollection> &run_data_collaction)
const {
279 if (run_data_collaction ==
nullptr) {
281 " run_data is null - cannot publish merged run data.");
284 if (gstreamer_run_map ==
nullptr) {
286 " no run streamer map available - run data will not be published.");
292 normalize_run_data(run_data_collaction);
294 for (
const auto &[name, gstreamer]: *gstreamer_run_map) {
295 if (gstreamer ==
nullptr) {
297 " null gstreamer instance for run streamer ", name);
300 gstreamer->publishRunData(run_data_collaction);
305void GRunAction::normalize_run_data(
const std::shared_ptr<GRunDataCollection> &run_data_collaction)
const {
306 if (run_data_collaction ==
nullptr) {
308 " run_data_collaction is null - cannot normalize run data.");
311 const int events_processed = run_data_collaction->get_events_processed();
312 if (events_processed <= 0) {
313 log->warning(FUNCTION_NAME,
314 " events_processed is ", events_processed,
315 " - skipping normalization.");
319 const double norm =
static_cast<double>(events_processed);
321 for (
auto &[sdName, dataCollection]: run_data_collaction->getMutableDataCollectionMap()) {
322 if (dataCollection ==
nullptr) {
323 log->warning(FUNCTION_NAME,
324 " detector ", sdName,
325 " has null data collection - skipping.");
329 auto &digitizedData = dataCollection->getMutableDigitizedData();
330 if (digitizedData.empty() || digitizedData.front() ==
nullptr) {
334 auto &digitized = digitizedData.front();
338 auto it = to_normalize.find(sdName);
339 if (it != to_normalize.end()) {
340 for (
const auto &varName: it->second) {
341 const auto intVars = digitized->getIntObservablesMap(0);
342 const auto intIt = intVars.find(varName);
343 if (intIt != intVars.end()) {
344 digitized->includeVariable(varName,
static_cast<double>(intIt->second) / norm);
348 const auto dblVars = digitized->getDblObservablesMap(0);
349 const auto dblIt = dblVars.find(varName);
350 if (dblIt != dblVars.end()) {
351 digitized->includeVariable(varName, dblIt->second / norm);
GBase(const std::shared_ptr< GOptions > &gopt, std::string logger_name="")
std::shared_ptr< GLogger > log
GRunAction(std::shared_ptr< GOptions > gopts, std::shared_ptr< gdynamicdigitization::dRoutinesMap > digi_map, std::shared_ptr< GAnalysisAccumulator > analysis_accumulator=nullptr, std::shared_ptr< GSROFactory > sro_factory=nullptr)
Constructs the run action.
Thread-local run object created by the GEMC run action.
Declares GRunAction, the run-lifecycle action for the GEMC actions module.
constexpr const char * GRUNACTION_LOGGER
Declares GRun, the thread-local run container used by the GEMC actions module.
std::uint64_t GSROEventId
Defines error codes used by the GEMC actions module.
std::shared_ptr< const gstreamersMap > gstreamersMapPtr(const std::shared_ptr< GOptions > &gopts, int thread_id=-1)
vector< GStreamerDefinition > getGStreamerDefinition(const std::shared_ptr< GOptions > &gopts)
constexpr int ERR_GDIGIMAP_NOT_EXISTING
constexpr int ERR_STREAMERMAP_NOT_EXISTING
constexpr int ERR_GRUNACTION_NOT_EXISTING
constexpr int ERR_PUBLISH_ERROR