-
Notifications
You must be signed in to change notification settings - Fork 6
/
Copy pathjobSubmitter.py
executable file
·599 lines (526 loc) · 23.2 KB
/
jobSubmitter.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
from __future__ import print_function
import os, subprocess, sys, stat, glob, shutil, tarfile, six
from optparse import OptionParser
from collections import defaultdict, OrderedDict
from datetime import datetime
from Condor.Production.parseConfig import list_callback, parser_dict
def date_convert(value):
return datetime.strptime(value, "%Y-%m-%d %H:%M:%S")
def date_callback(option, opt, value, parser):
if value is None: return
setattr(parser.values, option.dest, date_convert(value))
# minimal sed-like function
# patterns = [(in,out),(in,out)]
# currently doesn't handle regex
def pysed(lines,out,patterns):
with open(out,'w') as outfile:
for line in lines:
linetmp = line
for pattern,replace in six.iteritems(patterns):
linetmp = linetmp.replace(str(pattern),str(replace))
outfile.write(linetmp)
# run xrdfs ls using physical file name
def generalized_ls(redir, indir, minDate=None, maxDate=None):
checkDates = minDate is not None or maxDate is not None
extra = " -l --full-time" if checkDates else ""
# automatically split
if len(indir)==0:
redir, indir = split_pfn(redir)
cmd = None
if indir.startswith("/store/"):
if redir.startswith("root://"):
date_split_start = 1
cmd = "xrdfs {} ls {} {}".format(redir,extra.replace("--full-time",""),indir)
elif redir.startswith("gsiftp://"):
date_split_start = 5
cmd = "gfal-ls {}".format(redir+indir)
else:
raise ValueError("Unknown redir {}".format(redir))
else:
date_split_start = 5
cmd = "ls {} {}".format(extra,indir)
# todo: replace w/ XRootD python bindings?
popen_args = dict(
shell=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)
if cmd.startswith("xrdfs"):
# necessary to communicate w/ cmslpc at fnal
popen_args.update(dict(env=dict(os.environ,**{'XrdSecGSISRVNAMES': 'cmseos.fnal.gov'})))
if six.PY3:
popen_args.update(dict(encoding="utf-8"))
results = list(filter(
None,
subprocess.Popen(
cmd,
**popen_args
).communicate()[0].split('\n')
))
# for consistency with xrdfs ls, which includes full lfn path (/store/...) in output
if not cmd.startswith("xrdfs"):
results = [lfn+'/'+res for res in results]
if checkDates:
dates = [date_convert(' '.join(line.split()[date_split_start:date_split_start+2])) for line in results]
results = [line.split()[-1] for line,date in zip(results,dates) if (minDate is None or date>minDate) and (maxDate is None or date<maxDate)]
return results
def split_pfn(pfn):
if "://" in pfn:
psplit = pfn.find("/store")
lfn = pfn[psplit:]
xrd = pfn[:psplit]
else:
lfn = pfn
xrd = ""
return xrd, lfn
# backward compatibility
def pyxrdfsls(pfn, minDate=None, maxDate=None):
return generalized_ls(pfn, "", minDate, maxDate)
# run xrdcp
def pyxrdcp(a,b,verbose=True):
# xrootd dir
xrdcp = subprocess.Popen(
"xrdcp "+a+" "+b,
shell=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
encoding="utf-8",
env=dict(os.environ,**{'XrdSecGSISRVNAMES': 'cmseos.fnal.gov'})
)
xrdcp_result = xrdcp.communicate()
rc = xrdcp.returncode
if rc!=0 and verbose:
print("exit code "+str(rc)+", failure in xrdcp")
print(xrdcp_result[1])
return rc
# get regex capture group from file
def match(pattern,fname):
import re
with open(fname,'r') as file:
result = re.search(pattern,file.read())
if result is not None and len(result.groups())>0:
return result.group(1)
return None
class protoJob(object):
def __init__(self):
self.patterns = OrderedDict()
self.appends = []
self.queue = ""
self.njobs = 0
self.nums = []
self.jdl = ""
self.name = "job"
self.chainName = ""
def __repr__(self):
line = (
"protoJob:\n"
"\tname = "+str(self.name)+"\n"
"\tchainName = "+str(self.chainName)+"\n"
"\tnjobs = "+str(self.njobs)+"\n"
"\tjdl = "+str(self.jdl)+"\n"
"\tqueue = "+str(self.queue)+"\n"
"\tpatterns = "+str(self.patterns)+"\n"
"\tappends = "+str(self.appends)+"\n"
"\tnums = "+str(self.nums)
)
return line
def makeName(self,num):
return self.name+"_"+str(num)
class jobSubmitter(object):
def __init__(self,argv=None,parser=None):
if argv is None: argv = sys.argv[1:]
self.defaultStep1 = False
self.scripts = ["step1.sh","step2.sh"]
# dict of (string name, bool exclusive)
self.modes = {}
# define parser
if parser is None:
parser = OptionParser(add_help_option=False)
parser.add_option("--help", dest="help", action="store_true", default=False, help='show this help message')
# add sets of options
self.addDefaultOptions(parser)
self.addStep1Options(parser)
self.addExtraOptions(parser)
# parse & do help
(options, args) = parser.parse_args(args=argv)
if options.help:
parser.print_help()
sys.exit()
# check for option errors
self.checkDefaultOptions(options,parser)
self.checkStep1Options(options,parser)
self.checkExtraOptions(options,parser)
# set as members
for key in options.__dict__:
setattr(self,key,options.__dict__[key])
# other vars
self.protoJobs = []
self.jdlLines = []
self.njobs = 0
self.filesSet = set()
self.runSet = set()
self.missingNums = {}
self.missingLines = []
def run(self):
self.initRun()
# job generation
self.generateSubmission()
# loop over protojobs
for job in self.protoJobs:
self.runPerJob(job)
# final stuff
self.finishRun()
def initRun(self):
self.initStep1()
if self.missing:
self.initMissing()
elif self.clean:
self.initClean()
def runPerJob(self,job):
if self.prepare:
self.doPrepare(job)
# mutually exclusive
if self.count:
self.doCount(job)
elif self.submit:
self.doSubmit(job)
elif self.missing:
self.doMissing(job)
elif self.clean:
self.doClean(job)
def finishRun(self):
if self.count:
self.finishCount()
elif self.missing:
self.finishMissing()
elif self.clean:
self.finishClean()
def addDefaultOptions(self,parser):
# control options
parser.add_option("-c", "--count", dest="count", default=False, action="store_true", help="count the expected number of jobs (default = %default)")
parser.add_option("-p", "--prepare", dest="prepare", default=False, action="store_true", help="prepare job inputs and JDL files (default = %default)")
parser.add_option("-s", "--submit", dest="submit", default=False, action="store_true", help="submit jobs to condor (default = %default)")
parser.add_option("-m", "--missing", dest="missing", default=False, action="store_true", help="check for missing jobs (default = %default)")
parser.add_option("-r", "--resub", dest="resub", default="", help="make a resub script with specified name (default = %default)")
parser.add_option("--min-date", dest="minDate", type="string", default=None, action="callback", callback=date_callback, help="minimum date for files in missing mode (default = %default)")
parser.add_option("--max-date", dest="maxDate", type="string", default=None, action="callback", callback=date_callback, help="maximum date for files in missing mode (default = %default)")
parser.add_option("-l", "--clean", dest="clean", default=False, action="store_true", help="clean up log files (default = %default)")
parser.add_option("--clean-dir", dest="cleanDir", default=".", help="output dir for log file .tar.gz (default = %default)")
parser.add_option("-u", "--user", dest="user", default=parser_dict["common"]["user"], help="view jobs from this user (submitter) (default = %default)")
parser.add_option("-q", "--no-queue-arg", dest="noQueueArg", default=False, action="store_true", help="don't use -queue argument in condor_submit (default = %default)")
self.modes.update({
"count": 1,
"prepare": 0,
"submit": 1,
"missing": 1,
"clean": 1,
})
def checkDefaultOptions(self,options,parser):
nModes = 0
nExcls = 0
lModes = ""
lExcls = ""
for mode,excl in six.iteritems(self.modes):
lModes += mode + ", "
if excl: lExcls += mode + ", "
if getattr(options,mode):
nModes += 1
if excl: nExcls += 1
lModes = lModes[:-2]
lExcls = lExcls[:-2]
if nExcls>1:
parser.error("Modes "+lExcls+" are exclusive, pick one!")
if nModes==0:
parser.error("No operation mode selected! ("+lModes+")")
# if you use a different step1.sh, you might need to change these
def addStep1Options(self,parser):
self.defaultStep1 = True
parser.add_option("-k", "--keep", dest="keep", default=False, action="store_true", help="keep existing tarball for job submission (default = %default)")
parser.add_option("-n", "--no-voms", dest="novoms", default=False, action="store_true", help="skip check and use of voms proxy (default = %default)")
parser.add_option("-t", "--cmssw-method", dest="cmsswMethod", default="transfer", help="how to get CMSSW env: transfer, cmsrel, or address for xrdcp (default = %default)")
parser.add_option("--scram-workaround", dest="scramWorkaround", default=False, action="store_true", help="workaround for scram ProjectRename bug (default = %default)")
def checkStep1Options(self,options,parser):
if options.cmsswMethod!="transfer" and options.cmsswMethod!="cmsrel" and not options.cmsswMethod.startswith("root://"):
parser.error("Unknown CMSSW method: "+options.cmsswMethod)
# no need to retar if not submitting or not using tarball
if options.cmsswMethod=="cmsrel" or not options.submit:
options.keep = True
if options.novoms and options.cmsswMethod.startswith("root://"):
parser.error("Can't xrdcp CMSSW without voms proxy!")
def addExtraOptions(self,parser):
# job options
parser.add_option("--jdl", dest="jdl", default="jobExecCondor.jdl", help="JDL template file for job (default = %default)")
parser.add_option("--disk", dest="disk", default=1000000, help="specify amount of disk space per job [kB] (default = %default)")
parser.add_option("--memory", dest="memory", default=2000, help="specify amount of memory per job [MB] (default = %default)")
parser.add_option("--cpus", dest="cpus", default=1, help="specify number of CPU threads per job (default = %default)")
parser.add_option("--sites", dest="sites", default=parser_dict["submit"]["sites"], help="comma-separated list of sites for global pool running (default = %default)")
parser.add_option("--env", dest="env", default=None, help="args to run job in Singularity environment using cmssw-env (default = %default)")
parser.add_option("--intermediate", dest="intermediate", default=False, action="store_true", help="specify that this is an intermediate job in a chain to disable staging out (default = %default)")
parser.add_option("--singularity", dest="singularity", default="", help="specify singularity image for job (default = %default)")
def checkExtraOptions(self,options,parser):
pass
# in case you want to keep most but not all options from a section
def removeOptions(self,parser,*options):
for option in options:
if parser.has_option(option):
parser.remove_option(option)
def generateSubmission(self):
pass
def generatePerJob(self,job):
self.generateDefault(job)
self.generateStep1(job)
self.generateExtra(job)
self.generateJdl(job)
def initStep1(self):
# check for grid proxy and tarball
if self.defaultStep1:
cmd = "./checkVomsTar.sh"
if self.keep: cmd += " -k"
if self.novoms: cmd += " -n"
if not self.keep and self.cmsswMethod.startswith("root://"): cmd += " -i "+self.cmsswMethod
sp = subprocess.Popen(cmd, shell=True, stdin = sys.stdin, stdout = sys.stdout, stderr = sys.stderr)
sp.wait()
def initMissing(self):
# find finished jobs via output file list
self.filesSet = self.findFinished()
# find running jobs from condor
self.runSet = self.findRunning()
def initClean(self):
self.initMissing()
# subtract running jobs from finished jobs (in case resubmitted)
self.filesSet = self.filesSet - self.runSet
self.logdir = "logs"
if not os.path.isdir(self.logdir):
os.mkdir(self.logdir)
def generateDefault(self,job):
job.patterns["SCRIPTARGS"] = ",".join(self.scripts)
def generateStep1(self,job):
# command line args for step1
cmsswver = os.getenv("CMSSW_VERSION")
step1args = "-C "+cmsswver
if self.cmsswMethod == "cmsrel" or self.scramWorkaround:
# cmsrel needs scram arch
step1args += " -L "+os.getenv("SCRAM_ARCH")
if self.cmsswMethod.startswith("root://"):
# xrdcp needs input dir
step1args += " -X "+self.cmsswMethod
job.patterns["STEP1ARGS"] = step1args
if self.cmsswMethod=="transfer":
job.patterns["CMSSWVER"] = cmsswver
else:
job.patterns["CMSSWVER.tar.gz, "] = ""
if self.novoms:
job.patterns["x509userproxy = $ENV(X509_USER_PROXY)\n"] = ""
def generateExtra(self,job):
is_cms_connect = os.uname()[1]=="login.uscms.org" or os.uname()[1]=="login-el7.uscms.org"
os_version = match("[^0-9]*([0-9]+).*","/etc/redhat-release")
job.patterns.update([
("OSVERSION","rhel"+os_version),
("MYDISK",self.disk),
("MYMEMORY",self.memory),
("MYCPUS",self.cpus),
("ENVARGS",("-E {}".format(self.env) if self.env is not None else "")+(" -I" if self.intermediate else "")),
("SINGULARITYARGS",('+SingularityImage = "{}"'.format(self.singularity) if len(self.singularity)>0 else "")+("\nRequirements = HAS_SINGULARITY == True" if is_cms_connect else "")),
])
# special option for CMS Connect
if is_cms_connect:
if len(self.sites)>0: job.appends.append("+DESIRED_Sites = \""+self.sites+"\"")
job.appends.append("+AvoidSystemPeriodicRemove = True")
# special option for UMD
if "umd.edu" in os.uname()[1]:
if 'slc7' in os.environ['SCRAM_ARCH']:
job.appends.append("Requirements = (TARGET.OpSysMajorVer == 7)")
else:
job.appends.append("Requirements = (TARGET.OpSysMajorVer == 6)")
# left for the user: JOBNAME, EXTRAINPUTS, EXTRAARGS
def generateJdl(self,job):
job.jdl = self.jdl.replace(".jdl","_"+job.name+".jdl")
def doCount(self,job):
self.njobs += job.njobs
def finishCount(self):
print(str(self.njobs)+" jobs")
def doPrepare(self,job):
# get template contents (move into separate fn/store in self?)
if len(self.jdlLines)==0:
with open(self.jdl,'r') as jdlfile:
self.jdlLines = jdlfile.readlines()
# replace patterns
pysed(self.jdlLines,job.jdl,job.patterns)
# append appends & queue
with open(job.jdl,'a') as outfile:
for append_ in job.appends:
outfile.write(append_+"\n")
if self.noQueueArg: outfile.write(job.queue.replace("-queue","Queue")+"\n")
else: outfile.write("# "+job.queue.replace("-queue","Queue")+"\n")
def doSubmit(self,job):
if os.path.isfile(job.jdl):
if self.noQueueArg:
cmd = "condor_submit "+job.jdl
else:
queue = job.queue
# form should be: '-queue "..."'
if queue[7]!='"': queue = queue[:7]+'"'+queue[7:]
if queue[-1]!='"': queue = queue+'"'
cmd = "condor_submit "+job.jdl+" "+job.queue
os.system(cmd)
else:
print("Error: couldn't find "+job.jdl+", try running in prepare mode")
def editMissing(self,numlist,jdl,noQueueArg):
returnLines = []
if noQueueArg:
# get jdl lines for this job
with open(jdl,'r') as file:
jdlLines = [line for line in file]
# overwrite queue command in jdl
with open(jdl,'w') as file:
for line in jdlLines:
if line.startswith("Queue"):
file.write("#"+line)
file.write("Queue Process in "+','.join(map(str,numlist))+"\n")
else:
file.write(line)
returnLines.append('condor_submit '+jdl)
else:
returnLines.append('condor_submit '+jdl+' -queue "Process in '+','.join(map(str,numlist))+'"')
return returnLines
def doMissing(self,job):
jobSet, jobDict = self.findJobs(job)
# replace name if necessary
if len(job.chainName)>0:
runSetTmp = {x.replace(job.chainName,job.name) for x in self.runSet}
else:
runSetTmp = self.runSet
# find difference
diffSet = jobSet - self.filesSet - runSetTmp
diffList = list(sorted(diffSet))
numlist = sorted([jobDict[j] for j in diffList])
if len(diffList)>0:
if len(self.resub)>0:
self.missingLines.extend(self.editMissing(numlist,job.jdl,self.noQueueArg))
else:
self.missingLines.extend(diffList)
self.missingNums[job.jdl] = numlist
def findJobs(self,job):
jobSet = set()
jobDict = {}
for num in job.nums:
name = job.makeName(num)
jobSet.add(name)
jobDict[name] = num
return (jobSet, jobDict)
def finishMissing(self):
# provide results
if len(self.missingLines)>0:
if len(self.resub)>0:
self.makeResubmit(self.resub, self.missingLines)
else:
print('\n'.join(self.missingLines))
else:
print("No missing jobs!")
def finishedToJobName(self,val):
return val.split("/")[-1].replace(".root","")
def findFinished(self):
# find finished jobs via output file list
filesSet = set()
if hasattr(self,"output"):
files = generalized_ls(self.output,"",self.minDate,self.maxDate)
# basename
filesSet = set([ self.finishedToJobName(f) for f in files])
return filesSet
def tryToGetCondor(self):
# try to find condor bindings
if six.PY2:
for condorPath in glob.glob("/usr/lib64/python{}.*/site-packages".format(sys.version_info[0])):
if condorPath not in sys.path and os.path.isdir(condorPath):
sys.path.append(condorPath)
# try to import condor bindings
try:
global htcondor,classad
import htcondor,classad
except:
print('Could not import htcondor bindings!')
return False
return True
def runningToJobName(self,val):
return "_".join(val.replace(".stdout","").split('_')[:-1])
def findRunning(self):
runSet = set()
hasCondor = self.tryToGetCondor()
if not hasCondor:
print('"Missing jobs" check will not consider running jobs.')
return runSet
# exclude removed jobs
constraint = "JobStatus!=3"
if len(self.user)>0: constraint += ' && Owner=="'+self.user+'"'
for cname, collector in six.iteritems(parser_dict["collectors"]):
if cname not in parser_dict["schedds"]:
print("Error: no schedds provided for collector "+cname+", so it will be skipped.")
continue
else:
if len(collector)==0:
coll = htcondor.Collector()
else:
coll = htcondor.Collector(collector)
for sch in parser_dict["schedds"][cname].split(','):
try:
scheddAd = coll.locate(htcondor.DaemonTypes.Schedd, sch)
schedd = htcondor.Schedd(scheddAd)
for result in schedd.xquery(constraint,["Out"]):
runSet.add(self.runningToJobName(result["Out"]))
except:
print("Warning: could not locate schedd "+sch)
return runSet
def makeResubmit(self, resub, missingLines):
with open(resub,'w') as rfile:
rfile.write("#!/bin/bash\n\n")
for stmp in missingLines:
rfile.write(stmp+'\n')
# make executable
st = os.stat(rfile.name)
os.chmod(rfile.name, st.st_mode | stat.S_IXUSR | stat.S_IXGRP | stat.S_IXOTH)
def doClean(self,job):
jobSet, jobDict = self.findJobs(job)
finishedJobSet = jobSet & self.filesSet
# remove these jobs from global list
self.filesSet = self.filesSet - jobSet
for jobname in finishedJobSet:
# gets .condor, .stdout, .stderr
for fname in glob.glob(jobname+"_*.*"):
shutil.move(fname,self.logdir)
def finishClean(self):
# check if nothing to do
if len(os.listdir(self.logdir))==0:
# remove tmp dir
shutil.rmtree(self.logdir)
return
# create compressed tarfile
logname = "logs_tmp.tar.gz"
with tarfile.open(logname,"w:gz") as tar:
tar.add(self.logdir,arcname=self.logdir)
num_logs = 0
# check what is already in dir
if len(self.cleanDir)==0: self.cleanDir = "."
if self.cleanDir.startswith("root://"):
# xrootd dir
files = pyxrdfsls(self.cleanDir)
files = [f for f in files if f.endswith(".tar.gz")]
else:
# local dir
files = glob.glob(logname.replace("tmp","*"))
files = [f for f in files if f!=logname]
if len(files)>0:
num_logs = max([int(f.split("_")[-1].replace(".tar.gz","")) for f in files])+1
logname2 = logname.replace("tmp",str(num_logs))
rc = 1
# copy to dir
if self.cleanDir.startswith("root://"):
# xrootd dir
rc = pyxrdcp(logname,self.cleanDir+"/"+logname2)
else:
# local dir
shutil.copy2(logname,self.cleanDir+"/"+logname2)
rc = 0
if rc==0:
print("copied logs to "+self.cleanDir+"/"+logname2)
# remove tmp file
os.remove(logname)
# remove tmp dir
shutil.rmtree(self.logdir)