The Gaudi Framework
master (57817a18)
Toggle main menu visibility
Loading...
Searching...
No Matches
pTools.py
Go to the documentation of this file.
1
11
import
pickle
12
import
time
13
from
multiprocessing
import
Event
14
15
from
GaudiPython
import
FAILURE, SUCCESS, gbl
16
17
# Eoin Smith
18
# 3 Aug 2010
19
20
#
21
# This script contains the ancillary classes and functions used in the
22
# GaudiPython Parallel model
23
#
24
# Classes :
25
# - HistoAgent : In charge of extracting Histograms from their Transient Store
26
# on a reader/worker, communicating them to the writer, and on
27
# the writer receiving them and rebuilding a single store
28
#
29
# - FileRecordsAgent : Similar to HistoAgent, but for FileRecords Data
30
#
31
# - LumiFSR : FSR data from different workers needs to be carefully merged to
32
# replicate the serial version; this class aids in that task by
33
# representing an LHCb LumiFSR object as a python class
34
#
35
# - PackedCaloHypos : Pythonization of an LHCb class, used for inspecting some
36
# differences in serial and parallel output
37
#
38
# - Syncer : This class is responsible for syncing processes for a specified
39
# section of execution. For example, one Syncer object might be
40
# syncing Initialisation, one for Running, one for Finalisation.
41
# Syncer uses multiprocessing.Event() objects as flags which are
42
# visible across the N processes sharing them.
43
# IMPORTANT : The Syncer objects in the GaudiPython Parallel model
44
# ensure that there is no hanging; they in effect, allow a timeout
45
# for Initialisation, Run, Finalise on all processes
46
#
47
# - SyncMini : A class wrapper for a multiprocessing.Event() object
48
#
49
# Methods :
50
# - getEventNumber(evt) : pass a valid instance of the GaudiPython TES
51
# ( AppMgr().evtsvc() ) to this to get the current
52
# Event Number as an integer (even from RawEvents!)
53
#
54
55
# used to convert stored histos (in AIDA format) to ROOT format
56
aida2root = gbl.Gaudi.Utils.Aida2ROOT.aida2root
57
58
# =========================== Classes =========================================
59
60
61
class
HistoAgent
:
62
def
__init__
(self, gmpComponent):
63
self.
_gmpc
= gmpComponent
64
self.
hvt
= self.
_gmpc
.hvt
65
self.
histos
= []
66
self.
qin
= self.
_gmpc
.hq
67
self.
log
= self.
_gmpc
.log
68
69
# There are many methods for booking Histogram Objects to Histo store
70
# here they are collected in a dictionary, with key = a relevant name
71
self.
bookingDict
= {}
72
self.
bookingDict
[
"DataObject"
] = self.
bookDataObject
73
self.
bookingDict
[
"NTuple::Directory"
] = self.
bookDataObject
74
self.
bookingDict
[
"NTuple::File"
] = self.
bookDataObject
75
self.
bookingDict
[
"TH1D"
] = self.
bookTH1D
76
self.
bookingDict
[
"TH2D"
] = self.
bookTH2D
77
self.
bookingDict
[
"TH3D"
] = self.
bookTH3D
78
self.
bookingDict
[
"TProfile"
] = self.
bookTProfile
79
self.
bookingDict
[
"TProfile2D"
] = self.
bookTProfile2D
80
81
def
register
(self, tup):
82
# add a tuple of (worker-id, histoDict) to self.histos
83
assert
tup.__class__.__name__ ==
"tuple"
84
self.
histos
.append(tup)
85
86
def
Receive
(self):
87
hstatus = self.
_gmpc
.nWorkers + 1
# +1 for the Reader!
88
while
True
:
89
tup = self.
qin
.get()
90
if
tup ==
"HISTOS_SENT"
:
91
self.
log
.debug(
"received HISTOS_SENT message"
)
92
hstatus -= 1
93
if
not
hstatus:
94
break
95
else
:
96
self.
register
(tup)
97
self.
_gmpc
.sEvent.set()
98
self.
log
.info(
"Writer received all histo bundles and set sync event"
)
99
return
SUCCESS
100
101
def
RebuildHistoStore
(self):
102
"""
103
Rebuild the Histogram Store from the histos received by Receive()
104
If we have a histo which is not in the store,
105
book and fill it according to self.bookingDict
106
If we have a histo with a matching histo in the store,
107
add the two histos, remembering that aida2root must be used on
108
the Stored histo for compatibility.
109
"""
110
errors = 0
111
for
tup
in
self.
histos
:
112
workerID, histDict = tup
113
added = 0
114
booked = 0
115
116
for
n
in
histDict.keys():
117
o = histDict[n]
118
obj = self.
hvt
.retrieve(n)
119
120
if
obj:
121
try
:
122
aida2root
(obj).Add(o)
123
except
Exception:
124
self.
log
.warning(f
"FAILED TO ADD : {obj}"
)
125
errors += 1
126
added += 1
127
else
:
128
if
o.__class__.__name__
in
self.
bookingDict
.keys():
129
try
:
130
self.
bookingDict
[o.__class__.__name__](n, o)
131
except
Exception:
132
self.
log
.warning(
133
f
"FAILED TO REGISTER : {o.__class__.__name__}\tto{n}"
134
)
135
errors += 1
136
else
:
137
self.
log
.warning(
138
f
"No booking method for: {n}\t{type(o)}\t{o.__class__.__name__}"
139
)
140
errors += 1
141
booked += 1
142
hs = self.
hvt
.getHistoNames()
143
self.
log
.info(
"Histo Store Rebuilt : "
)
144
self.
log
.info(f
" Contains {len(hs)} objects."
)
145
self.
log
.info(f
" Errors in Rebuilding : {errors}"
)
146
return
SUCCESS
147
148
def
bookDataObject
(self, n, o):
149
"""
150
Register a DataObject to the Histo Store
151
"""
152
self.
_gmpc
.hvt.registerObject(n, o)
153
154
def
bookTH1D
(self, n, o):
155
"""
156
Register a ROOT 1D THisto to the Histo Store
157
"""
158
obj = self.
hvt
._ihs.book(
159
n,
160
o.GetTitle(),
161
o.GetXaxis().GetNbins(),
162
o.GetXaxis().GetXmin(),
163
o.GetXaxis().GetXmax(),
164
)
165
aida2root
(obj).Add(o)
166
167
def
bookTH2D
(self, n, o):
168
"""
169
Register a ROOT 2D THisto to the Histo Store
170
"""
171
obj = self.
hvt
._ihs.book(
172
n,
173
o.GetTitle(),
174
o.GetXaxis().GetNbins(),
175
o.GetXaxis().GetXmin(),
176
o.GetXaxis().GetXmax(),
177
o.GetYaxis().GetNbins(),
178
o.GetYaxis().GetXmin(),
179
o.GetYaxis().GetXmax(),
180
)
181
aida2root
(obj).Add(o)
182
183
def
bookTH3D
(self, n, o):
184
"""
185
Register a ROOT 3D THisto to the Histo Store
186
"""
187
obj = self.
hvt
._ihs.book(
188
n,
189
o.GetTitle(),
190
o.GetXaxis().GetXbins(),
191
o.GetXaxis().GetXmin(),
192
o.GetXaxis().GetXmax(),
193
o.GetYaxis().GetXbins(),
194
o.GetYaxis().GetXmin(),
195
o.GetYaxis().GetXmax(),
196
o.GetZaxis().GetXbins(),
197
o.GetZaxis().GetXmin(),
198
o.GetZaxis().GetXmax(),
199
)
200
aida2root
(obj).Add(o)
201
202
def
bookTProfile
(self, n, o):
203
"""
204
Register a ROOT TProfile to the Histo Store
205
"""
206
obj = self.
hvt
._ihs.bookProf(
207
n,
208
o.GetTitle(),
209
o.GetXaxis().GetNbins(),
210
o.GetXaxis().GetXmin(),
211
o.GetXaxis().GetXmax(),
212
o.GetOption(),
213
)
214
aida2root
(obj).Add(o)
215
216
def
bookTProfile2D
(self, n, o):
217
"""
218
Register a ROOT TProfile2D to the Histo Store
219
"""
220
obj = self.
hvt
._ihs.bookProf(
221
n,
222
o.GetTitle(),
223
o.GetXaxis().GetNbins(),
224
o.GetXaxis().GetXmin(),
225
o.GetXaxis().GetXmax(),
226
o.GetYaxis().GetNbins(),
227
o.GetYaxis().GetXmin(),
228
o.GetYaxis().GetXmax(),
229
)
230
aida2root
(obj).Add(o)
231
232
233
# =============================================================================
234
235
236
class
FileRecordsAgent
:
237
def
__init__
(self, gmpComponent):
238
self.
_gmpc
= gmpComponent
239
self.
fsr
= self.
_gmpc
.fsr
240
self.
q
= self.
_gmpc
.fq
241
self.
log
= self.
_gmpc
.log
242
self.
objectsIn
= []
# used for collecting FSR store objects
243
self.
objectsOut
= []
244
245
def
localCmp
(self, tupA, tupB):
246
# sort tuples by a particular element
247
# for the sort() method
248
ind = 0
249
valA = tupA[ind]
250
valB = tupB[ind]
251
if
valA < valB:
252
return
-1
253
elif
valA > valB:
254
return
1
255
else
:
256
return
0
257
258
def
SendFileRecords
(self):
259
# send the FileRecords data as part of finalisation
260
261
# Take Care of FileRecords!
262
# There are two main things to consider here
263
# 1) The DataObjects in the FileRecords Transient Store
264
# 2) The fact that they are Keyed Containers, containing other objects
265
#
266
# The Lead Worker, nodeID=0, sends everything in the FSR store, as
267
# a completeness guarantee,
268
#
269
# send in form ( nodeID, path, object)
270
self.
log
.info(
"Sending FileRecords..."
)
271
lst = self.
fsr
.getHistoNames()
272
273
# Check Validity
274
if
not
lst:
275
self.
log
.info(
"No FileRecords Data to send to Writer."
)
276
self.
q
.put(
"END_FSR"
)
277
return
SUCCESS
278
279
# no need to send the root node
280
if
"/FileRecords"
in
lst:
281
lst.remove(
"/FileRecords"
)
282
283
for
l
in
lst:
284
o = self.
fsr
.retrieveObject(l)
285
if
hasattr(o,
"configureDirectAccess"
):
286
o.configureDirectAccess()
287
# lead worker sends everything, as completeness guarantee
288
if
self.
_gmpc
.nodeID == 0:
289
self.
objectsOut
.append((0, l, pickle.dumps(o)))
290
else
:
291
# only add the Event Counter
292
# and non-Empty Keyed Containers (ignore empty ones)
293
if
l ==
"/FileRecords/EventCountFSR"
:
294
tup = (self.
_gmpc
.nodeID, l, pickle.dumps(o))
295
self.
objectsOut
.append(tup)
296
elif
"KeyedContainer"
in
o.__class__.__name__:
297
# It's a Keyed Container
298
nObjects = o.numberOfObjects()
299
if
nObjects:
300
self.
log
.debug(f
"Keyed Container {l} with {nObjects} objects"
)
301
tup = (self.
_gmpc
.nodeID, l, pickle.dumps(o))
302
self.
objectsOut
.append(tup)
303
else
:
304
self.
log
.info(f
"Ignoring {o.__class__.__name__} in FSR"
)
305
306
self.
log
.debug(
"Done with FSR store, just to send to Writer."
)
307
308
if
self.
objectsOut
:
309
self.
log
.debug(f
"{len(self.objectsOut)} FSR objects to Writer"
)
310
for
ob
in
self.
objectsOut
:
311
self.
log
.debug(f
"\t{ob[0]}"
)
312
self.
q
.put(self.
objectsOut
)
313
else
:
314
self.
log
.info(
"Valid FSR Store, but no data to send to Writer"
)
315
self.
log
.info(
"SendFSR complete"
)
316
self.
q
.put(
"END_FSR"
)
317
return
SUCCESS
318
319
def
Receive
(self):
320
# Receive contents of all Workers FileRecords Transient Stores
321
self.
log
.info(
"Receiving FSR store data..."
)
322
nc = self.
_gmpc
.nWorkers
323
while
nc > 0:
324
objects = self.
q
.get()
325
if
objects ==
"END_FSR"
:
326
nc -= 1
327
continue
328
if
nc == 0:
329
break
330
# but if it's regular objects...
331
for
o
in
objects:
332
self.
objectsIn
.append(o)
333
# Now sort it by which worker it came from
334
# an object is : (nodeID, path, pickledObject)
335
self.
objectsIn
.sort(cmp=self.
localCmp
)
336
self.
log
.info(
"All FSR data received"
)
337
return
SUCCESS
338
339
def
Rebuild
(self):
340
# objects is a list of (path, serializedObject) tuples
341
for
sourceNode, path, serialob
in
self.
objectsIn
:
342
self.
log
.debug(f
"Working with {path}"
)
343
ob = pickle.loads(serialob)
344
if
hasattr(ob,
"update"
):
345
ob.update()
346
if
hasattr(ob,
"numberOfObjects"
):
347
nCont = ob.numberOfObjects()
348
self.
log
.debug(f
"\t {type(ob).__name__} has containedObjects : {nCont}"
)
349
if
sourceNode == 0:
350
self.
log
.debug(f
"Registering Object to : {path}"
)
351
self.
fsr
.registerObject(path, ob)
352
else
:
353
self.
log
.debug(f
"Merging Object to : {path}"
)
354
self.
MergeFSRobject
(sourceNode, path, ob)
355
# As RecordStream has been split into Worker and Writer parts, the
356
# count for output is wrong... fix that here, as every event received
357
# by the writer is written (validation testing occurs on the worker)
358
359
self.
log
.info(
"FSR Store Rebuilt. Correcting EventCountFSR"
)
360
if
bool(self.
fsr
._idp):
# There might not be an FSR stream (Gauss)
361
ecount =
"/FileRecords/EventCountFSR"
362
if
self.
fsr
[ecount]:
363
self.
fsr
[ecount].setOutput(self.
_gmpc
.nIn)
364
self.
log
.info(
365
f
"Event Counter Output set : {ecount} : {self.fsr[ecount].output()}"
366
)
367
# Do some reporting
368
self.
log
.debug(
"FSR store reconstructed!"
)
369
lst = self.
fsr
.getHistoNames()
370
if
lst:
371
for
l
in
lst:
372
ob = self.
fsr
.retrieveObject(l)
373
if
hasattr(ob,
"configureDirectAccess"
):
374
ob.configureDirectAccess()
375
if
hasattr(ob,
"containedObjects"
):
376
# if ob.numberOfObjects() :
377
self.
log
.debug(
378
f
"\t{l} (cont. objects : {ob.numberOfObjects()})"
379
)
380
else
:
381
self.
log
.debug(f
"\t{l}"
)
382
self.
log
.info(
"FSR Store fully rebuilt."
)
383
return
SUCCESS
384
385
def
MergeFSRobject
(self, sourceNode, path, ob):
386
# Merge Non-Empty Keyed Container from Worker>0
387
if
path ==
"/FileRecords/TimeSpanFSR"
:
388
# TimeSpanFSR is a straightforward case
389
self.
ProcessTimeSpanFSR
(path, ob)
390
elif
path ==
"/FileRecords/EventCountFSR"
:
391
# Event Counter is also easy
392
self.
ProcessEventCountFSR
(path, ob)
393
# now other cases may not be so easy...
394
elif
(
395
"KeyedContainer"
in
ob.__class__.__name__
396
and
"LumiFSR"
in
ob.__class__.__name__
397
):
398
# Keyed Container of LumiFSRs : extract and re-register
399
self.
MergeLumiFSR
(path, ob)
400
else
:
401
self.
log
.info(f
"Skipping Merge of {ob.__class__.__name__} at {path}"
)
402
403
def
ProcessTimeSpanFSR
(self, path, ob):
404
ob2 = self.
fsr
.retrieveObject(path)
405
if
ob.containedObjects().size():
406
sz = ob.containedObjects().size()
407
cob = ob2.containedObjects()[0]
408
min = cob.earliest()
409
max = cob.latest()
410
for
j
in
range(sz):
411
cob = ob.containedObjects()[j]
412
self.
log
.debug(
"Adding TimeSpanFSR"
)
413
if
cob.earliest() < min:
414
min = cob.earliest()
415
if
cob.latest() > max:
416
max = cob.latest()
417
# this is annoying: it has to be rebuilt, without a key & added
418
continue
419
tsfsr = gbl.LHCb.TimeSpanFSR()
420
tsfsr.setEarliest(min)
421
tsfsr.setLatest(max)
422
self.
fsr
[path].clear()
423
self.
fsr
[path].add(tsfsr)
424
425
def
ProcessEventCountFSR
(self, path, ob):
426
self.
log
.debug(
"Event Count Input Addition"
)
427
self.
fsr
[path].setInput(self.
fsr
[path].input() + ob.input())
428
429
def
MergeLumiFSR
(self, path, keyedC):
430
from
ROOT
import
string
431
432
# Fetch the first lumi
433
keyedContainer = self.
fsr
.retrieveObject(path)
434
# The LumiFSR KeyedContainers only have one object
435
assert
keyedContainer.numberOfObjects() == 1
436
l = keyedContainer.containedObject(0)
437
baseLumi =
LumiFSR
(l)
438
# Now deal with the argument Non-empty Keyed Container of LumiFSRs
439
nCont = keyedC.numberOfObjects()
440
for
i
in
range(nCont):
441
obj = keyedC.containedObject(i)
442
nextLumi =
LumiFSR
(obj)
443
baseLumi.merge(nextLumi)
444
# Now Rebuild and ReRegister
445
newLumi = gbl.LHCb.LumiFSR()
446
for
r
in
baseLumi.runs:
447
newLumi.addRunNumber(r)
448
for
f
in
baseLumi.files:
449
newLumi.addFileID(
string
(f))
450
for
k
in
baseLumi.keys:
451
increment, integral = baseLumi.info[k]
452
newLumi.addInfo(k, increment, integral)
453
# clear existing Keyed Container
454
self.
fsr
[path].clear()
455
# Add newly merged lumiFSR
456
self.
fsr
[path].add(newLumi)
457
return
SUCCESS
458
459
460
# =============================================================================
461
462
463
class
LumiFSR
:
464
def
__init__
(self, lumi):
465
# lumi looks like :
466
# { runs : 69857 69858
467
# files : root:/castor/cer.../069857_0000000006.raw
468
# info (key/incr/integral) : 0 8 0 / 1 8 259 / 2 8 76 ... }
469
470
# class variables
471
self.
runs
= []
472
self.
files
= []
473
self.
info
= {}
474
self.
keys
= []
475
476
# get run numbers
477
for
r
in
lumi.runNumbers():
478
self.
runs
.append(r)
479
# get file ids
480
for
f
in
lumi.fileIDs():
481
self.
files
.append(f)
482
# Now the tricky bit, the info is not accessible via Python
483
# except as a string
484
s = str(lumi)
485
sa = s.split(
"info (key/incr/integral) : "
)[-1]
486
sa = sa.split(
"/"
)[:-1]
487
for
rec
in
sa:
488
k, i, t = rec.split()
489
k = int(k)
490
i = int(i)
491
t = int(t)
492
self.
info
[k] = (i, t)
493
self.
keys
= self.
info
.
keys
()
494
495
def
merge
(self, otherLumi):
496
assert
otherLumi.__class__.__name__ ==
"LumiFSR"
497
# add any extra runs
498
for
r
in
otherLumi.runs:
499
if
r
in
self.
runs
:
500
pass
501
else
:
502
self.
runs
.append(r)
503
self.
runs
.sort()
504
# add any extra fileIDs
505
for
f
in
otherLumi.files:
506
if
f
in
self.
files
:
507
pass
508
else
:
509
self.
files
.append(f)
510
self.
files
.sort()
511
# Now add any extra records
512
for
k
in
otherLumi.keys:
513
increment, integral = otherLumi.info[k]
514
if
k
in
self.
keys
:
515
myIncrement, myIntegral = self.
info
[k]
516
self.
info
[k] = (myIncrement + increment, myIntegral + integral)
517
else
:
518
self.
info
[k] = (increment, integral)
519
# don't forget to update keys
520
self.
keys
= self.
info
.
keys
()
521
522
def
__repr__
(self):
523
s =
"LumiFSR Python class\n"
524
s +=
"\tRuns : \n"
525
for
r
in
self.
runs
:
526
s += f
"\t\t{r}\n"
527
s +=
"\tFiles : \n"
528
for
f
in
self.
files
:
529
s += f
"\t\t{f}\n"
530
s +=
"\tInfo : \n"
531
for
k
in
self.
keys
:
532
increment, integral = self.
info
[k]
533
s += f
"\t\t{k}\t{increment}\t{integral}\n"
534
return
s
535
536
537
# =============================================================================
538
539
540
class
PackedCaloHypo
:
541
def
__init__
(self, o):
542
cl =
"LHCb::PackedCaloHypo"
543
assert
o.__class__.__name__ == cl
544
self.
centX
= o.centX
545
self.
centY
= o.centY
546
self.
cerr
= (o.cerr00, o.cerr10, o.cerr11)
547
self.
cov
= (o.cov00, o.cov10, o.cov11, o.cov20, o.cov21, o.cov22)
548
self.
firstCluster
= o.firstCluster
549
self.
firstDigit
= o.firstDigit
550
self.
firstHypo
= o.firstHypo
551
self.
hypothesis
= o.hypothesis
552
self.
key
= o.key
553
self.
lastCluster
= o.lastCluster
554
self.
lastDigit
= o.lastDigit
555
self.
lastHypo
= o.lastHypo
556
self.
lh
= o.lh
557
self.
pos
= (o.posE, o.posX, o.posY)
558
self.
z
= o.z
559
560
def
__repr__
(self):
561
s =
"PackedCaloHypo : \n"
562
s += f
"\tcentX : {self.centX}\n"
563
s += f
"\tcentY : {self.centY}\n"
564
s += f
"\tcerr : {self.cerr}\n"
565
s += f
"\tcov : {self.cov}\n"
566
s += f
"\tfirstCluster : {self.firstCluster}\n"
567
s += f
"\tfirstDigit : {self.firstDigit}\n"
568
s += f
"\tfirstHypo : {self.firstHypo}\n"
569
s += f
"\thypothesis : {self.hypothesis}\n"
570
s += f
"\tkey : {self.key}\n"
571
s += f
"\tlastCluster : {self.lastCluster}\n"
572
s += f
"\tlastDigit : {self.lastDigit}\n"
573
s += f
"\tlastHypo : {self.lastHypo}\n"
574
s += f
"\tlh : {self.lh}\n"
575
s += f
"\tpos : {self.pos}\n"
576
s += f
"\tz : {self.z}\n"
577
s +=
"---------------------------------------\n"
578
return
s
579
580
581
# =============================================================================
582
583
584
class
SyncMini
:
585
def
__init__
(self, event, lastEvent=None):
586
self.
event
= event
587
self.
t
= 0.0
588
self.
lastEvent
=
None
589
if
lastEvent:
590
self.
lastEvent
= lastEvent
591
592
def
check
(self):
593
return
self.
event
.is_set()
594
595
def
checkLast
(self):
596
return
self.
lastEvent
.is_set()
597
598
def
reset
(self):
599
self.
event
.clear()
600
self.
t
= time.time()
601
602
def
getTime
(self):
603
return
self.
t
604
605
def
set
(self):
606
self.
event
.
set
()
607
608
def
__repr__
(self):
609
s =
"---------- SyncMini --------------\n"
610
s += f
" Status : {self.event.is_set()}\n"
611
s += f
" t : {self.t:5.2f}\n"
612
if
self.
lastEvent
:
613
s += f
"Last Event : {self.lastEvent.is_set()}\n"
614
s +=
"----------------------------------\n"
615
return
s
616
617
618
# =============================================================================
619
620
621
class
Syncer
:
622
def
__init__
(
623
self, nWorkers, log, manyEvents=False, limit=None, step=None, firstEvent=None
624
):
625
# Class to help synchronise the sub-processes
626
self.
limit
= limit
627
self.
step
= step
628
self.
d
= {}
629
self.
manyEvents
= manyEvents
630
631
for
i
in
range(-2, nWorkers):
632
self.
d
[i] =
SyncMini
(Event(), lastEvent=Event())
633
if
self.
manyEvents
:
634
self.
limitFirst
= firstEvent
635
636
self.
keys
= list(self.
d
.
keys
())
637
self.
nWorkers
= nWorkers
638
self.
log
= log
639
640
def
syncAll
(self, step="Not specified"):
641
# is it this method, or is it the rolling version needed?
642
# if so, drop through...
643
644
if
self.
manyEvents
:
645
sc = self.
syncAllRolling
()
646
return
sc
647
648
# Regular version ----------------------------
649
for
i
in
range(0, self.
limit
, self.
step
):
650
if
self.
checkAll
():
651
self.
log
.info(f
"{step} : All procs done @ {i} s"
)
652
break
653
else
:
654
time.sleep(self.
step
)
655
656
# Now the time limit is up... check the status one final time
657
if
self.
checkAll
():
658
self.
log
.info(f
"All processes : {step} ok."
)
659
return
SUCCESS
660
else
:
661
self.
log
.critical(f
"Some process is hanging on : {step}"
)
662
for
k
in
self.
keys
:
663
self.
log
.critical(f
"{step} : Proc/Stat : {k}/{self.d[k].check()}"
)
664
return
FAILURE
665
666
def
syncAllRolling
(self):
667
# Keep track of the progress of Event processing
668
# Each process syncs after each event, so keep clearing
669
# the sync Event, and re-checking
670
# Note the time between True checks too, if the time
671
# between events exceeds singleEvent, this is considered a hang
672
673
# set the initial time
674
begin = time.time()
675
firstEv = {}
676
timers = {}
677
for
k
in
self.
keys
:
678
self.
d
[k].reset()
679
firstEv[k] =
False
680
timers[k] = 0.0
681
682
active = self.
keys
683
while
True
:
684
# check the status of each sync object
685
for
k
in
active:
686
sMini = self.
d
[k]
687
688
if
sMini.check()
or
sMini.checkLast():
689
if
sMini.checkLast()
and
sMini.check():
690
# if last Event set,then event loop finished
691
active.remove(k)
692
alive = time.time() - begin
693
self.
log
.info(f
"Audit : Node {k} alive for {alive:5.2}"
)
694
else
:
695
sMini.reset()
696
else
:
697
# the event still has not been checked, how long is that?
698
# is it the first Event?
699
wait = time.time() - sMini.getTime()
700
cond = wait > self.
limit
701
if
not
firstEv[k]:
702
cond = wait > self.
limitFirst
703
firstEv[k] =
True
704
if
cond:
705
# It is hanging!
706
self.
log
.critical(f
"Single event wait : {wait:5.2f}"
)
707
self.
processHang
()
708
return
FAILURE
709
710
# Termination Criteria : if all procs have been removed, we're done
711
if
self.
checkLastEvents
():
712
self.
log
.info(
"TC met for event loop"
)
713
break
714
else
:
715
# sleep, loop again
716
time.sleep(self.
step
)
717
718
self.
log
.info(
"All processes Completed all Events ok"
)
719
return
SUCCESS
720
721
def
processHang
(self):
722
self.
log
.critical(
"Some proc is hanging during Event processing!"
)
723
for
k
in
self.
keys
:
724
self.
log
.critical(f
"Proc/Stat : {k} / {self.d[k].check()}"
)
725
return
726
727
def
checkAll
(self):
728
# Check the status of each Sync object
729
# return True or False
730
currentStatus = [mini.check()
for
mini
in
self.
d
.values()]
731
return
all(currentStatus)
732
733
def
checkLastEvents
(self):
734
# check if all of the lastEvents are set to true in self.d[k][1]
735
stat = [sMini.checkLast()
for
sMini
in
self.
d
.values()]
736
return
all(stat)
737
738
739
# =========================== Methods =========================================
740
741
742
def
getEventNumber
(evt):
743
# The class-independent version of the Event Number Retrieval method
744
#
745
n =
None
746
# First Attempt : Unpacked Event Data
747
lst = [
"/Event/Gen/Header"
,
"/Event/Rec/Header"
]
748
for
l
in
lst:
749
try
:
750
n = evt[l].evtNumber()
751
return
n
752
except
Exception:
753
# No evt number at this path
754
continue
755
756
# second attepmt : try DAQ/RawEvent data
757
# The Evt Number is in bank type 16, bank 0, data pt 4
758
try
:
759
n = evt[
"/Event/DAQ/RawEvent"
].banks(16)[0].data()[4]
760
return
n
761
except
Exception:
762
pass
763
764
# Default Action
765
return
n
766
767
768
# ================================= EOF =======================================
GaudiMP.pTools.FileRecordsAgent
Definition
pTools.py:236
GaudiMP.pTools.FileRecordsAgent._gmpc
_gmpc
Definition
pTools.py:238
GaudiMP.pTools.FileRecordsAgent.Rebuild
Rebuild(self)
Definition
pTools.py:339
GaudiMP.pTools.FileRecordsAgent.__init__
__init__(self, gmpComponent)
Definition
pTools.py:237
GaudiMP.pTools.FileRecordsAgent.SendFileRecords
SendFileRecords(self)
Definition
pTools.py:258
GaudiMP.pTools.FileRecordsAgent.MergeLumiFSR
MergeLumiFSR(self, path, keyedC)
Definition
pTools.py:429
GaudiMP.pTools.FileRecordsAgent.MergeFSRobject
MergeFSRobject(self, sourceNode, path, ob)
Definition
pTools.py:385
GaudiMP.pTools.FileRecordsAgent.Receive
Receive(self)
Definition
pTools.py:319
GaudiMP.pTools.FileRecordsAgent.ProcessTimeSpanFSR
ProcessTimeSpanFSR(self, path, ob)
Definition
pTools.py:403
GaudiMP.pTools.FileRecordsAgent.objectsOut
list objectsOut
Definition
pTools.py:243
GaudiMP.pTools.FileRecordsAgent.log
log
Definition
pTools.py:241
GaudiMP.pTools.FileRecordsAgent.fsr
fsr
Definition
pTools.py:239
GaudiMP.pTools.FileRecordsAgent.q
q
Definition
pTools.py:240
GaudiMP.pTools.FileRecordsAgent.ProcessEventCountFSR
ProcessEventCountFSR(self, path, ob)
Definition
pTools.py:425
GaudiMP.pTools.FileRecordsAgent.localCmp
localCmp(self, tupA, tupB)
Definition
pTools.py:245
GaudiMP.pTools.FileRecordsAgent.objectsIn
list objectsIn
Definition
pTools.py:242
GaudiMP.pTools.HistoAgent
Definition
pTools.py:61
GaudiMP.pTools.HistoAgent.bookTProfile2D
bookTProfile2D(self, n, o)
Definition
pTools.py:216
GaudiMP.pTools.HistoAgent.register
register(self, tup)
Definition
pTools.py:81
GaudiMP.pTools.HistoAgent.qin
qin
Definition
pTools.py:66
GaudiMP.pTools.HistoAgent.bookDataObject
bookDataObject(self, n, o)
Definition
pTools.py:148
GaudiMP.pTools.HistoAgent.__init__
__init__(self, gmpComponent)
Definition
pTools.py:62
GaudiMP.pTools.HistoAgent.hvt
hvt
Definition
pTools.py:64
GaudiMP.pTools.HistoAgent.bookTH1D
bookTH1D(self, n, o)
Definition
pTools.py:154
GaudiMP.pTools.HistoAgent._gmpc
_gmpc
Definition
pTools.py:63
GaudiMP.pTools.HistoAgent.bookTH2D
bookTH2D(self, n, o)
Definition
pTools.py:167
GaudiMP.pTools.HistoAgent.RebuildHistoStore
RebuildHistoStore(self)
Definition
pTools.py:101
GaudiMP.pTools.HistoAgent.histos
list histos
Definition
pTools.py:65
GaudiMP.pTools.HistoAgent.bookingDict
dict bookingDict
Definition
pTools.py:71
GaudiMP.pTools.HistoAgent.log
log
Definition
pTools.py:67
GaudiMP.pTools.HistoAgent.Receive
Receive(self)
Definition
pTools.py:86
GaudiMP.pTools.HistoAgent.bookTProfile
bookTProfile(self, n, o)
Definition
pTools.py:202
GaudiMP.pTools.HistoAgent.bookTH3D
bookTH3D(self, n, o)
Definition
pTools.py:183
GaudiMP.pTools.LumiFSR
Definition
pTools.py:463
GaudiMP.pTools.LumiFSR.keys
list keys
Definition
pTools.py:474
GaudiMP.pTools.LumiFSR.__repr__
__repr__(self)
Definition
pTools.py:522
GaudiMP.pTools.LumiFSR.merge
merge(self, otherLumi)
Definition
pTools.py:495
GaudiMP.pTools.LumiFSR.__init__
__init__(self, lumi)
Definition
pTools.py:464
GaudiMP.pTools.LumiFSR.info
dict info
Definition
pTools.py:473
GaudiMP.pTools.LumiFSR.runs
list runs
Definition
pTools.py:471
GaudiMP.pTools.LumiFSR.files
list files
Definition
pTools.py:472
GaudiMP.pTools.PackedCaloHypo
Definition
pTools.py:540
GaudiMP.pTools.PackedCaloHypo.cov
tuple cov
Definition
pTools.py:547
GaudiMP.pTools.PackedCaloHypo.__repr__
__repr__(self)
Definition
pTools.py:560
GaudiMP.pTools.PackedCaloHypo.key
key
Definition
pTools.py:552
GaudiMP.pTools.PackedCaloHypo.lastCluster
lastCluster
Definition
pTools.py:553
GaudiMP.pTools.PackedCaloHypo.z
z
Definition
pTools.py:558
GaudiMP.pTools.PackedCaloHypo.__init__
__init__(self, o)
Definition
pTools.py:541
GaudiMP.pTools.PackedCaloHypo.lastHypo
lastHypo
Definition
pTools.py:555
GaudiMP.pTools.PackedCaloHypo.firstHypo
firstHypo
Definition
pTools.py:550
GaudiMP.pTools.PackedCaloHypo.centX
centX
Definition
pTools.py:544
GaudiMP.pTools.PackedCaloHypo.centY
centY
Definition
pTools.py:545
GaudiMP.pTools.PackedCaloHypo.lastDigit
lastDigit
Definition
pTools.py:554
GaudiMP.pTools.PackedCaloHypo.firstCluster
firstCluster
Definition
pTools.py:548
GaudiMP.pTools.PackedCaloHypo.cerr
tuple cerr
Definition
pTools.py:546
GaudiMP.pTools.PackedCaloHypo.firstDigit
firstDigit
Definition
pTools.py:549
GaudiMP.pTools.PackedCaloHypo.hypothesis
hypothesis
Definition
pTools.py:551
GaudiMP.pTools.PackedCaloHypo.pos
tuple pos
Definition
pTools.py:557
GaudiMP.pTools.PackedCaloHypo.lh
lh
Definition
pTools.py:556
GaudiMP.pTools.SyncMini
Definition
pTools.py:584
GaudiMP.pTools.SyncMini.getTime
getTime(self)
Definition
pTools.py:602
GaudiMP.pTools.SyncMini.reset
reset(self)
Definition
pTools.py:598
GaudiMP.pTools.SyncMini.lastEvent
lastEvent
Definition
pTools.py:588
GaudiMP.pTools.SyncMini.event
event
Definition
pTools.py:586
GaudiMP.pTools.SyncMini.__init__
__init__(self, event, lastEvent=None)
Definition
pTools.py:585
GaudiMP.pTools.SyncMini.checkLast
checkLast(self)
Definition
pTools.py:595
GaudiMP.pTools.SyncMini.check
check(self)
Definition
pTools.py:592
GaudiMP.pTools.SyncMini.set
set(self)
Definition
pTools.py:605
GaudiMP.pTools.SyncMini.t
float t
Definition
pTools.py:587
GaudiMP.pTools.SyncMini.__repr__
__repr__(self)
Definition
pTools.py:608
GaudiMP.pTools.Syncer
Definition
pTools.py:621
GaudiMP.pTools.Syncer.keys
keys
Definition
pTools.py:636
GaudiMP.pTools.Syncer.d
dict d
Definition
pTools.py:628
GaudiMP.pTools.Syncer.checkAll
checkAll(self)
Definition
pTools.py:727
GaudiMP.pTools.Syncer.limitFirst
limitFirst
Definition
pTools.py:634
GaudiMP.pTools.Syncer.step
step
Definition
pTools.py:627
GaudiMP.pTools.Syncer.checkLastEvents
checkLastEvents(self)
Definition
pTools.py:733
GaudiMP.pTools.Syncer.nWorkers
nWorkers
Definition
pTools.py:637
GaudiMP.pTools.Syncer.manyEvents
manyEvents
Definition
pTools.py:629
GaudiMP.pTools.Syncer.limit
limit
Definition
pTools.py:626
GaudiMP.pTools.Syncer.syncAllRolling
syncAllRolling(self)
Definition
pTools.py:666
GaudiMP.pTools.Syncer.processHang
processHang(self)
Definition
pTools.py:721
GaudiMP.pTools.Syncer.log
log
Definition
pTools.py:638
GaudiMP.pTools.Syncer.syncAll
syncAll(self, step="Not specified")
Definition
pTools.py:640
GaudiMP.pTools.Syncer.__init__
__init__(self, nWorkers, log, manyEvents=False, limit=None, step=None, firstEvent=None)
Definition
pTools.py:624
string
STL class.
GaudiMP.pTools.aida2root
aida2root
Definition
pTools.py:56
GaudiMP.pTools.getEventNumber
getEventNumber(evt)
Definition
pTools.py:742
GaudiMP
python
GaudiMP
pTools.py
Generated on
for The Gaudi Framework by
1.17.0