actions
Loading...
Searching...
No Matches
gRunAction.cc
Go to the documentation of this file.
1// gemc
2#include "gRunAction.h"
3#include "gRun.h"
5#include "gutsConventions.h"
6
7// geant4
8#include "G4Threading.hh"
9
10std::mutex GRunAction::completed_run_data_mutex;
11GRunAction::CompletedRunData GRunAction::completed_worker_run_data;
12
13
14// Construct the run action and retain access to shared configuration and
15// digitization services for the current execution context.
16GRunAction::GRunAction(std::shared_ptr<GOptions> gopt,
17 std::shared_ptr<gdynamicdigitization::dRoutinesMap> digi_map,
18 std::shared_ptr<GAnalysisAccumulator> analyzer,
19 std::shared_ptr<GSROFactory> sro) : GBase(gopt, GRUNACTION_LOGGER),
20 goptions(std::move(gopt)),
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());
25 log->debug(CONSTRUCTOR, FUNCTION_NAME, desc);
26}
27
28
29// Create the thread-local run object used by Geant4 for this execution context.
30G4Run *GRunAction::GenerateRun() {
31 log->debug(NORMAL, FUNCTION_NAME);
32
33 return new GRun(goptions, digitization_routines_map);
34}
35
36// Initialize run-scoped bookkeeping, determine which streamer categories are
37// needed for this run, and open the appropriate connections.
38void GRunAction::BeginOfRunAction(const G4Run *aRun) {
39 const auto thread_id = G4Threading::G4GetThreadId();
40 const auto run = aRun->GetRunID();
41
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>();
47 }
48
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) {
53 log->error(gstreamer::ERR_PUBLISH_ERROR, "SRO begin run: ", error.what());
54 }
55 }
56
57 // Reset the per-run mode flags before scanning the digitization routines.
58 need_a_thread_streamer = false;
59 need_a_run_streamer = false;
60
61 // Inspect the available digitization routines to determine whether this run
62 // requires event-mode publication, run-mode publication, or both.
63 if (digitization_routines_map != nullptr) {
64 for (const auto &[plugin, digiRoutine]: *digitization_routines_map) {
65 if (digiRoutine == nullptr) {
66 log->error(gaction::ERR_GDIGIMAP_NOT_EXISTING, FUNCTION_NAME,
67 " null digitization routine registered for plugin ", plugin);
68 }
69
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;
75 }
76 }
77 } else {
78 log->error(gaction::ERR_GDIGIMAP_NOT_EXISTING, FUNCTION_NAME,
79 " digitization_routines_map is null - streamer mode detection skipped.");
80 }
81
82 for (const auto& streamer_definition : gstreamer::getGStreamerDefinition(goptions)) {
83 if (streamer_definition.type == "event") {
84 need_a_thread_streamer = true;
85 break;
86 }
87 }
88
89
90 // Worker threads own event-mode publication, so they lazily build and open
91 // thread-local streamers only when at least one event-mode digitizer exists.
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);
95 gstreamer_threads_map = gstreamer::gstreamersMapPtr(goptions, thread_id);
96 }
97
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.");
101 }
102
103 for (const auto &[name, gstreamer]: *gstreamer_threads_map) {
104 if (gstreamer == nullptr) {
106 "Null GStreamer entry ", name, " in thread ", thread_id);
107 }
108
109 if (!gstreamer->openConnection()) {
111 "Failed to open connection for GStreamer ", name,
112 " in thread ", thread_id);
113 }
114
115 log->info(2, FUNCTION_NAME, "Worker thread [", thread_id, "]: opening connection for ",
116 guts::KGRN, name, guts::RST,
117 " for run ", run, ". Number of events to be processed: ", neventsThisRun);
118 }
119 }
120 // The master thread owns run-mode publication, so it opens the run streamers
121 // only when at least one digitizer accumulates payload at run scope.
122 else if (IsMaster() && need_a_run_streamer) {
123 if (gstreamer_run_map == nullptr) {
124 log->info(1, "Defining run gstreamers for run ", run);
125 gstreamer_run_map = gstreamer::gstreamersMapPtr(goptions);
126 }
127
128 if (gstreamer_run_map == nullptr) {
129 log->error(1, FUNCTION_NAME, " gstreamer_run_map is null in master thread ",
130 " - cannot open connections.");
131 }
132
133 for (const auto &[name, gstreamer]: *gstreamer_run_map) {
134 if (gstreamer == nullptr) {
136 "Null GStreamer entry ", name, " in master thread");
137 }
138
139 if (!gstreamer->openConnection()) {
141 "Failed to open connection for GStreamer in master thread ", name);
142 }
143
144 log->info(2, FUNCTION_NAME, "Master Thread: opening connection for ",
145 guts::KGRN, name, guts::RST,
146 " for run ", run, ". Number of events to be processed: ", neventsThisRun);
147 }
148 }
149}
150
151// Close streamers at run end and, when running on the master thread, merge and
152// publish the run-level payload accumulated by workers.
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";
157 // In MT Geant4 reaches the master here after worker run actions; sequential mode owns both roles.
158 if (sro_factory && IsMaster()) {
159 try { sro_factory->finish_run(aRun->GetNumberOfEvent() == aRun->GetNumberOfEventToBeProcessed()); }
160 catch (const std::exception& error) {
161 // finish_run has already joined all crate/progress threads before reporting the first error.
162 log->error(gstreamer::ERR_PUBLISH_ERROR, "SRO finish run: ", error.what());
163 }
164 }
165
166 if (!IsMaster() && need_a_thread_streamer) {
167 if (gstreamer_threads_map == nullptr) {
168 log->error(gaction::ERR_STREAMERMAP_NOT_EXISTING, FUNCTION_NAME,
169 " gstreamer_map is null in thread ", thread_id,
170 " - cannot close connections.");
171 } else {
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);
175
176 if (gstreamer == nullptr) {
178 "Null GStreamer entry ", name, " in thread ", thread_id);
179 }
180
181 if (!gstreamer->closeConnection()) {
182 log->error(1, "Failed to close connection for GStreamer ", name, " in thread ", thread_id);
183 }
184 }
185 }
186 }
187
188 // Each execution context writes only to its private shard. The shared accumulator locks once here.
189 if (analysis_accumulator != nullptr && analysis_shard != nullptr) {
190 if (!analysis_shard->empty()) { analysis_accumulator->merge(std::move(*analysis_shard)); }
191 analysis_shard.reset();
192 }
193
194 // Worker threads do not publish merged run data. Instead, they hand their
195 // completed run-level accumulation to the shared pool and return.
196 if (!IsMaster()) {
197 // Only contribute to the master merge pool when a run-mode streamer will drain it;
198 // otherwise the pool is never taken and would accumulate stale prior-run data.
199 if (need_a_run_streamer) { stash_worker_run_data(); }
200 return;
201 }
202
203 if (IsMaster() && need_a_run_streamer) {
204 // Gather all worker-produced run data for this run and merge them into a
205 // single master-side run-data object before publication.
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);
210
211 std::shared_ptr<GRunDataCollection> merged_run_data;
212
213 for (auto &worker_run_data: completed_run_data) {
214 if (worker_run_data == nullptr) {
215 continue;
216 }
217
218 // Create the merged destination lazily only if there is at least one
219 // valid worker contribution to merge.
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));
223 }
224
225 merged_run_data->merge(*worker_run_data);
226 }
227
228 // Publish the merged run-level payload once, after all workers have contributed.
229 if (merged_run_data != nullptr) {
230 publish_run_data(merged_run_data);
231 }
232
233 if (gstreamer_run_map == nullptr) {
234 log->error(gaction::ERR_STREAMERMAP_NOT_EXISTING, FUNCTION_NAME,
235 " gstreamer_map is null in master thread - cannot close connections.");
236 }
237
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);
241
242 if (gstreamer == nullptr) {
244 "Null GStreamer entry ", name, " in master thread");
245 }
246
247 if (!gstreamer->closeConnection()) {
248 log->error(1, "Failed to close connection for GStreamer ", name, " in master thread");
249 }
250 }
251 }
252
253}
254
255// Move this worker's completed run-data object into the protected static pool
256// so it can later be collected by the master thread.
257void GRunAction::stash_worker_run_data() {
258 if (run_data == nullptr) {
259 return;
260 }
261
262 std::scoped_lock lock(completed_run_data_mutex);
263 completed_worker_run_data.emplace_back(std::move(run_data));
264}
265
266// Extract and clear the protected pool of completed worker run-data objects.
267auto GRunAction::take_completed_worker_run_data() -> CompletedRunData {
268 std::scoped_lock lock(completed_run_data_mutex);
269
270 auto result = std::move(completed_worker_run_data);
271 completed_worker_run_data.clear();
272
273 return result;
274}
275
276
277// Publish the merged run-level payload to every configured master-side run streamer.
278void GRunAction::publish_run_data(const std::shared_ptr<GRunDataCollection> &run_data_collaction) const {
279 if (run_data_collaction == nullptr) {
280 log->error(gaction::ERR_GRUNACTION_NOT_EXISTING, FUNCTION_NAME,
281 " run_data is null - cannot publish merged run data.");
282 }
283
284 if (gstreamer_run_map == nullptr) {
285 log->error(gaction::ERR_STREAMERMAP_NOT_EXISTING, FUNCTION_NAME,
286 " no run streamer map available - run data will not be published.");
287 }
288
289 // Normalize once, before publishing: normalize_run_data() mutates run_data_collaction
290 // in place and is NOT idempotent, so running it per streamer would divide the run
291 // observables by events_processed once for every configured run streamer.
292 normalize_run_data(run_data_collaction);
293
294 for (const auto &[name, gstreamer]: *gstreamer_run_map) {
295 if (gstreamer == nullptr) {
296 log->error(gaction::ERR_STREAMERMAP_NOT_EXISTING, FUNCTION_NAME,
297 " null gstreamer instance for run streamer ", name);
298 }
299
300 gstreamer->publishRunData(run_data_collaction);
301 }
302}
303
304
305void GRunAction::normalize_run_data(const std::shared_ptr<GRunDataCollection> &run_data_collaction) const {
306 if (run_data_collaction == nullptr) {
307 log->error(gaction::ERR_GRUNACTION_NOT_EXISTING, FUNCTION_NAME,
308 " run_data_collaction is null - cannot normalize run data.");
309 }
310
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.");
316 return;
317 }
318
319 const double norm = static_cast<double>(events_processed);
320
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.");
326 continue;
327 }
328
329 auto &digitizedData = dataCollection->getMutableDigitizedData();
330 if (digitizedData.empty() || digitizedData.front() == nullptr) {
331 continue;
332 }
333
334 auto &digitized = digitizedData.front();
335 // we are in a const method so we can't loop directky over to_normalize[sdName]
336 // cause this call could modify the map!
337
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);
345 continue;
346 }
347
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);
352 }
353 }
354 }
355 }
356}
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.
Definition gRunAction.cc:16
Thread-local run object created by the GEMC run action.
Definition gRun.h:54
Declares GRunAction, the run-lifecycle action for the GEMC actions module.
constexpr const char * GRUNACTION_LOGGER
Definition gRunAction.h:29
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.
run
#define FUNCTION_NAME
CONSTRUCTOR
NORMAL
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
constexpr char RST[]
constexpr char KGRN[]