8from pathlib
import Path
10from AnaAlgorithm.DualUseConfig
import isAthena
11from AnaAlgorithm.Logging
import logging
13logCPGridRun = logging.getLogger(
'CPGridRun')
39 from AnalysisAlgorithmsConfig.AthenaCPRunScript
import AthenaCPRunScript
42 from AnalysisAlgorithmsConfig.EventLoopCPRunScript
import EventLoopCPRunScript
47 parser = argparse.ArgumentParser(description=
'CPGrid runscript to submit CPRun.py jobs to the grid. '
48 'This script will submit a job to the grid using files in the input text one by one.'
49 'CPRun.py can handle multiple sources of input and create one output; but not this script',
51 formatter_class=argparse.RawTextHelpFormatter)
52 parser.add_argument(
'-h',
'--help', dest=
'help', action=
'store_true', help=
'Show this help message and continue')
54 ioGroup = parser.add_argument_group(
'Input/Output file configuration')
55 ioGroup.add_argument(
'-i',
'--input-list', dest=
'input_list', help=
'Path to the text file containing list of containers on the panda grid. Each container will be passed to prun as --inDS and is run individually')
56 ioGroup.add_argument(
'--output-files', dest=
'output_files', nargs=
'+', default=[
'output.root'],
57 help=
'The output files of the grid job. Example: --output-files A.root B.txt B.root results in A/A.root, B/B.txt, B/B.root in the output directory. No need to specify if using CPRun.py')
58 ioGroup.add_argument(
'--destSE', dest=
'destSE', default=
'', type=str, help=
'Destination storage element (PanDA)')
59 ioGroup.add_argument(
'--mergeType', dest=
'mergeType', default=
'Default', type=str, help=
'Output merging type, [None, Default, xAOD]')
61 pandaGroup = parser.add_argument_group(
'Input/Output naming configuration')
62 pandaGroup.add_argument(
'--gridUsername', dest=
'gridUsername', default=os.getenv(
'USER',
''), type=str, help=
'Grid username, or the groupname. Default is the current user. Only affect file naming')
63 pandaGroup.add_argument(
'--prefix', dest=
'prefix', default=
'', type=str, help=
'Prefix for the output directory. Dynamically set with input container if not provided')
64 pandaGroup.add_argument(
'--suffix', dest=
'suffix', default=
'',type=str, help=
'Suffix for the output directory')
65 pandaGroup.add_argument(
'--outDS', dest=
'outDS', default=
'', type=str,
66 help=
'Name of an output dataset. outDS will contain all output files (PanDA). If not provided, support dynamic naming if input name is in the Atlas production format or typical user production format')
68 cpgridGroup = parser.add_argument_group(
'CPGrid configuration')
69 cpgridGroup.add_argument(
'--groupProduction', dest=
'groupProduction', action=
'store_true', help=
'Only use for official production')
71 cpgridGroup.add_argument(
'--exec', dest=
'exec', type=str,
72 help=
'Executable line for the CPRun.py or custom script to run on the grid encapsulated in a double quote (PanDA)\n'
73 'Run CPRun.py with preset behavior including streamlined file i/o. E.g, "CPRun.py -t config.yaml --no-systematics".\n'
74 'Run custom script: "customRun.py -i inputs -o output --text-config config.yaml --flagA --flagB"\n'
77 submissionGroup = parser.add_argument_group(
'Submission configuration')
78 submissionGroup.add_argument(
'--noSubmit', dest=
'noSubmit', action=
'store_true', help=
'Do not submit the job to the grid (PanDA). Useful to inspect the prun command')
79 submissionGroup.add_argument(
'--testRun', dest=
'testRun', action=
'store_true', help=
'Will submit job to the grid but greatly limit the number of files per job (10) and number of events (300)')
80 submissionGroup.add_argument(
'--recreateTar', dest=
'recreateTar', action=
'store_true', help=
'Re-compress the source code. Source code are compressed by default in submission, this is useful when the source code is updated')
81 submissionGroup.add_argument(
'--useCentralPackage', dest=
'useCentralPackage', action=
'store_true', help=
'Use central package instead of custom packages')
82 submissionGroup.add_argument(
'--bulk-submission', dest=
'bulk_submission', action=
'store_true', help=
'Submit all containers in the input list as one task.')
84 miscGroup = parser.add_argument_group(
'Miscellaneous configuration')
85 miscGroup.add_argument(
'-y',
'--agreeAll', dest=
'agreeAll', action=
'store_true', help=
'Agree to all the submission details without asking for confirmation. Use with caution!')
86 miscGroup.add_argument(
'--checkInputDS', dest=
'checkInputDS', action=
'store_true', help=
'Check if the input datasets are available on the AMI.')
87 miscGroup.add_argument(
'--framework', dest=
'framework', default=
'CPGridRun', type=str, help=
'Declaring a name for your submission for PanDA team to collect statistics. Default is CPGridRun')
95 converting unknown args to a dictionary
98 if unknownArgsDict
and self.
hasPrun():
100 logCPGridRun.info(f
"Adding prun exclusive arguments: {unknownArgsDict.keys()}")
101 elif unknownArgsDict:
102 logCPGridRun.warning(f
"Unknown arguments detected: {unknownArgsDict}. Cannot check the availability in Prun because Prun is not available / noSubmit is on.")
105 return unknownArgsDict
109 if not self.
args.input_list:
110 raise ValueError(
'No input list provided, use --input-list to specify the input containers')
112 input_list_path = Path(self.
args.input_list)
113 if input_list_path.exists():
114 if input_list_path.suffix ==
'.txt':
116 elif input_list_path.suffix ==
'.json':
117 raise NotImplementedError(
'JSON input list parsing is not implemented')
119 raise ValueError(
'Unsupported input list format, only .txt files are supported.')
120 elif CPGridRun.isAtlasProductionFormat(self.
args.input_list):
124 raise ValueError(
'use --input-list to specify input containers')
129 for output
in self.
args.output_files:
131 output_files.extend(output.split(
','))
133 output_files.append(output)
138 logCPGridRun.info(
"\033[92m\n If you are using CPRun.py, the following flags are for the CPRun.py in this framework\033[0m")
139 self.
_runscript.parser.usage = argparse.SUPPRESS
150 self.
cmd[input] = cmd
151 self.
outputs[input] = config[
"outDS"]
158 'cmtConfig': os.environ[
"CMTCONFIG"],
159 'writeInputToTxt':
'IN:in.txt',
163 'addNthFieldOfInDSToLFN':
'2,3,6',
164 'framework': self.
args.framework,
166 if self.
args.noSubmit:
167 config[
'noSubmit'] =
True
169 if self.
args.mergeType ==
'xAOD':
170 config[
'mergeScript'] =
'xAODMerge %OUT `echo %IN | sed \'s/,/ /g\'`'
172 if self.
args.mergeType !=
'None':
173 config[
'mergeOutput'] =
True
176 if self.
args.useCentralPackage:
178 config[
'noBuild'] =
True
179 config[
'noCompile'] =
True
180 config[
'athenaTag'] = f
"AnalysisBase,{os.environ['AnalysisBase_VERSION']}"
182 config[
'outTarBall'] = self.
_tarfile
183 config[
'useAthenaPackages'] =
True
187 config[
'useAthenaPackages'] =
True
189 if self.
args.groupProduction:
190 config[
'official'] =
True
191 config[
'voms'] = f
'atlas:/atlas/{self.args.gridUsername}/Role=production'
194 config[
'destSE'] = self.
args.destSE
196 if self.
args.testRun:
197 config[
'nEventsPerFile'] = 100
201 for k, v
in config.items():
202 if isinstance(v, bool)
and v:
204 elif v
is not None and v !=
'':
206 value = v
if k ==
'exec' else shlex.quote(str(v))
207 cmd += f
'--{k} {value} \\\n'
208 return cmd.rstrip(
' \\\n'), config
212 Cleans the unknown args by removing leading dashes and ensuring they are in key-value pairs
214 unknown_args_dict = {}
222 unknown_args_dict[self.
unknown_args[idx].lstrip(
'-')] =
True
224 return unknown_args_dict
228 check the arguments against the prun script to ensure they are valid
229 See https://github.com/PanDAWMS/panda-client/blob/master/pandaclient/PrunScript.py
231 import pandaclient.PrunScript
233 original_argv = sys.argv
236 prunArgsDict = pandaclient.PrunScript.main(get_options=
True)
237 sys.argv = original_argv
238 nonPrunOrCPGridArgs = []
240 if arg
not in prunArgsDict:
241 nonPrunOrCPGridArgs.append(arg)
242 if nonPrunOrCPGridArgs:
243 logCPGridRun.error(f
"Unknown arguments detected: {nonPrunOrCPGridArgs}. They do not belong to CPGridRun or Panda.")
244 raise ValueError(f
"Unknown arguments detected: {nonPrunOrCPGridArgs}. They do not belong to CPGridRun or Panda.")
247 for key, cmd
in self.
cmd.items():
248 parsed_name = CPGridRun.atlasProductionNameParser(key)
249 logCPGridRun.info(
"\n"
251 "\n".join([f
" {k.replace('_', ' ').title()}: {v}" for k, v
in parsed_name.items()]))
252 logCPGridRun.info(f
"Command: \n{cmd}")
260 import pyAMI.atlas.api
261 except ModuleNotFoundError:
263 "Cannot import pyAMI, please run the following commands:\n\n"
266 "voms-proxy-init -voms atlas\n"
268 "and make sure you have a valid certificate.")
276 client = pyAMI.client.Client(
'atlas')
277 pyAMI.atlas.api.init()
281 results = pyAMI.atlas.api.list_datasets(client, patterns=queries)
282 except pyAMI.exception.Error:
284 "Cannot query AMI, please run 'voms-proxy-init -voms atlas' and ensure your certificate is valid.")
291 Helper function to prepare a list of queries for the AMI based on the input list.
292 It will replace the _p### with _p% to match the latest ptag.
295 regex = re.compile(
"_p[0-9]+")
298 for datasetName
in self.
cmd:
299 parsed = CPGridRun.atlasProductionNameParser(datasetName)
300 datasetPtag[datasetName] = parsed.get(
'ptag')
301 queries.append(regex.sub(
"_p%", datasetName))
302 return queries, datasetPtag
306 regex = re.compile(
"_p[0-9]+")
307 results = [r[
'ldn']
for r
in results]
311 for datasetName
in self.
cmd:
312 if datasetName
not in results:
313 notFound.append(datasetName)
315 base = regex.sub(
"_p%", datasetName)
316 matching = [r
for r
in results
if r.startswith(base.replace(
"_p%",
""))]
318 mParsed = CPGridRun.atlasProductionNameParser(m)
320 mPtagInt = int(mParsed.get(
'ptag',
'p0')[1:])
321 currentPtagInt = int(datasetPtag.get(datasetName,
'p0')[1:])
322 if mPtagInt > currentPtagInt:
323 latestPtag[datasetName] = f
"p{mPtagInt}"
324 except (ValueError, TypeError):
328 logCPGridRun.info(
"Newer version of datasets found in AMI:")
329 for name, ptag
in latestPtag.items():
330 logCPGridRun.info(f
"{name} -> ptag: {ptag}")
333 logCPGridRun.error(
"Some input datasets are not available in AMI, missing datasets are likely to fail on the grid:")
334 logCPGridRun.error(
", ".join(notFound))
346 if CPGridRun.isAtlasProductionFormat(name)
and not label:
353 {group/user}.{username}.{prefix}.{DSID}.{format}.{tags}.{suffix}
355 nameParser = CPGridRun.atlasProductionNameParser(name)
356 base =
'group' if self.
args.groupProduction
else 'user'
357 username = self.
args.gridUsername
358 dsid = nameParser[
'DSID']
359 tags =
'_'.join(nameParser[
'tags'])
360 fileFormat = nameParser[
'format']
361 prefix = self.
args.prefix
if self.
args.prefix
else nameParser[
'main'].
split(
'_')[0]
364 result = [base, username, prefix, dsid, fileFormat, tags, suffix]
365 return ".".join(filter(
None, result))
369 {group/user}.{username}.{prefix}.{main}.{suffix}
371 parts = name.split(
'.')
372 base =
'group' if self.
args.groupProduction
else 'user'
373 username = self.
args.gridUsername
374 main = label
if label
else parts[2]
375 main = f
"{self.args.prefix}.{main}" if self.
args.prefix
else main
376 main = f
"{main}.{self.args.suffix}" if self.
args.suffix
else main
378 result = [base, username, main]
379 return ".".join(filter(
None, result))
383 return self.
args.suffix
384 if self.
args.testRun:
386 return f
"test_{uuid.uuid4().hex[:6]}"
391 tarball_mtime = os.path.getmtime(self.
_tarfile)
if os.path.exists(self.
_tarfile)
else 0
396 for root, _, files
in os.walk(buildDir):
398 file_path = os.path.join(root, file)
400 if os.path.getmtime(file_path) > tarball_mtime:
401 logCPGridRun.info(f
"File {file_path} is newer than the tarball.")
403 except FileNotFoundError:
407 if sourceDir
is None:
408 logCPGridRun.warning(
"Source directory is not detected, auto-compression is not performed. Use --recreateTar to update the submission")
410 for root, _, files
in os.walk(sourceDir):
412 file_path = os.path.join(root, file)
414 if os.path.getmtime(file_path) > tarball_mtime:
415 logCPGridRun.info(f
"File {file_path} is newer than the tarball.")
417 except FileNotFoundError:
422 buildDir = os.environ[
"CMAKE_PREFIX_PATH"]
423 buildDir = os.path.dirname(buildDir.split(
":")[0])
427 cmakeCachePath = os.path.join(self.
_buildDir(),
'CMakeCache.txt')
429 if not os.path.exists(cmakeCachePath):
431 with open(cmakeCachePath,
'r')
as cmakeCache:
432 for line
in cmakeCache:
433 if '_SOURCE_DIR:STATIC=' in line:
434 sourceDir = line.split(
'=')[1].
strip()
439 if not self.
args.exec:
440 raise ValueError(
'No exec command provided, use --exec to specify the command to run on the grid')
443 isCPRunDefault = self.
args.exec.startswith(
'-')
or self.
args.exec.startswith(
'CPRun.py')
445 'input_list':
'in.txt',
446 'merge_output_files': len(self.
args.output_files) == 1,
448 if not isCPRunDefault:
449 if self.
_isFirstRun: logCPGridRun.warning(
"Non-CPRun.py is detected, please ensure the exec string is formatted correctly. Exec string will not be automatically formatted.")
450 return f
'"{self.args.exec}"'
454 runscriptArgs, unknownArgs = self.
_runscript.parser.parse_known_args(self.
args.exec.split(
' '))
457 unknown_flags = [arg
for arg
in unknownArgs
if arg.startswith(
'--')]
459 logCPGridRun.error(f
"Unknown flags detected in the exec string: {unknown_flags}. Please check the exec string.")
460 raise ValueError(f
"Unknown arguments detected: {unknown_flags}")
463 for key, value
in formatingClause.items():
464 if hasattr(runscriptArgs, key):
465 old_value = getattr(runscriptArgs, key)
466 if old_value
is None or old_value == self.
_runscript.parser.get_default(key):
467 setattr(runscriptArgs, key, value)
468 if self.
_isFirstRun: logCPGridRun.info(f
"Setting '{key}' to '{value}' (CPRun.py default is: '{old_value}')")
470 if self.
_isFirstRun: logCPGridRun.warning(f
"Preserving user-defined '{key}': '{old_value}', default formatting '{value}' will not be applied.")
472 logCPGridRun.error(f
"Formatting clause '{key}' is not recognized in the CPRun.py script. Check CPGridRun.py")
473 raise ValueError(f
"Formatting clause '{key}' is not recognized in the CPRun.py script. Check CPGridRun.py")
476 arg_string =
' '.join(
477 f
'--{k.replace("_", "-")}' if isinstance(v, bool)
and v
else
478 f
'--{k.replace("_", "-")} {v}' for k, v
in vars(runscriptArgs).items()
if v
not in [
None,
False]
480 return f
'"CPRun.py {arg_string}"'
483 from AnalysisAlgorithmsConfig.CPBaseRunner
import CPBaseRunner
484 if not hasattr(runscriptArgs,
'text_config'):
485 self.
_errorCollector[
'no yaml'] =
"No YAML configuration file is specified in the exec string. Please provide one using --text-config"
487 yamlPath = getattr(runscriptArgs,
'text_config')
489 haveLocalYaml = CPBaseRunner.findLocalPathYamlConfig(yamlPath)
491 logCPGridRun.warning(
"A path to a local YAML configuration file is found, but it may not be grid-usable.")
493 repoYamls, _ = CPBaseRunner.findRepoPathYamlConfig(yamlPath)
494 if repoYamls
and len(repoYamls) > 1:
495 self.
_errorCollector[
'ambiguous yamls'] = f
'Multiple files named \"{yamlPath}\" found in the analysis repository. Please provide a more specific path to the config file.\nMatches found:\n' +
'\n'.join(repoYamls)
497 elif repoYamls
and len(repoYamls) == 1:
498 logCPGridRun.info(f
"Found a grid-usable YAML configuration file in the analysis repository: {repoYamls[0]}")
501 if haveLocalYaml
and self.
args.useCentralPackage:
502 logCPGridRun.warning(
"A path to a local YAML configuration file is found, no custom packages are found, proceed with /cvmfs packages only.")
504 if not repoYamls
and not self.
args.useCentralPackage:
505 self.
_errorCollector[
'no usable yaml'] = f
"Grid usable YAML configuration file not found: {yamlPath}"
507 self.
_errorCollector[
'have local yaml'] = f
"Only a local YAML configuration file is found: {yamlPath}, not usable in the grid.\n" \
508 f
"Make sure the YAML file is in build/x86_64-el9-gcc14-opt/data/package_name/config.yaml. You can install the YAML file through CMakeList.txt with `atlas_install_data( data/* )`; use `-t package_name/config.yaml` in the --exec\n"\
509 f
"Or if you are only using central packages, please use the `--useCentralPackage` flag."
512 outputs = [f
'{output.split(".")[0]}:{output}' if ":" not in output
else output
for output
in self.
args.output_files]
513 return ','.join(outputs)
517 prun_path = shutil.which(
"prun")
518 if prun_path
is None:
520 "The 'prun' command is not found. If you are on lxplus, please run the following commands:\n\n"
523 "voms-proxy-init -voms atlas\n"
525 "Make sure you have a valid certificate."
532 for key, cmd
in self.
cmd.items():
533 logCPGridRun.info(f
"Submitting: {self.outputs[key]}")
534 process = subprocess.Popen(cmd, shell=
True, stdout=sys.stdout, stderr=sys.stderr)
535 process.communicate()
540 name = name.split(
":")[1]
542 if name.startswith(
"mc")
or name.startswith(
"data"):
545 logCPGridRun.warning(
"Name is not in the Atlas production format, assuming it is a user production")
551 The custom name has many variations, but most of them follow user/group.username.datasetname.suffix
554 parts = filename.split(
'.')
555 result[
'userType'] = parts[0]
556 result[
'username'] = parts[1]
557 result[
'main'] = parts[2]
558 result[
'suffix'] = parts[-1]
564 Parsing file name into a dictionary, an example is given here
565 mc20_13TeV.410470.PhPy8EG_A14_ttbar_hdamp258p75_nonallhad.deriv.DAOD_PHYS.e6337_s3681_r13167_p5855/DAOD_PHYS.34865530._000740.pool.root.1
567 datasetName: mc20_13TeV.410470.PhPy8EG_A14_ttbar_hdamp258p75_nonallhad.deriv.DAOD_PHYS.e6337_s3681_r13167_p5855
568 projectName: mc20_13TeV
572 main: PhPy8EG_A14_ttbar_hdamp258p75_nonallhad
573 TODO generator: PhPy8Eg
574 TODO tune: A14 # For Pythia8
576 TODO hdamp: 258p75 # For Powheg
577 TODO decayType: nonallhad
580 tags: e###_s###_r###_p###_a###_t###_b#
581 etag: e6337 # EVNT (EVGEN) production and merging
582 stag: s3681 # Geant4 simulation to produce HITS and merging!
583 rtag: r13167 # Digitisation and reconstruction, as well as AOD merging
584 ptag: p5855 # Production of NTUP_PILEUP format and merging
585 atag: aXXX: atlfast configuration (both simulation and digit/recon)
586 ttag: tXXX: tag production configuration
587 btag: bXXX: bytestream production configuration
600 datasetPart, filePart = filename.split(
'/')
602 datasetPart = filename
606 if ':' in datasetPart:
607 datasetPart = datasetPart.split(
':')[1]
610 if datasetPart.startswith(
'user')
or datasetPart.startswith(
'group'):
611 result[
'datasetName'] = datasetPart
615 datasetParts = datasetPart.split(
'.')
616 result[
'datasetName'] = datasetPart
618 result[
'projectName'] = datasetParts[0]
620 campaign_energy = result[
'projectName'].
split(
'_')
621 result[
'campaign'] = campaign_energy[0]
622 result[
'energy'] = campaign_energy[1]
625 result[
'DSID'] = datasetParts[1]
626 result[
'main'] = datasetParts[2]
627 result[
'step'] = datasetParts[3]
628 result[
'format'] = datasetParts[4]
631 tags = datasetParts[5].
split(
'_')
632 result[
'tags'] = tags
634 if tag.startswith(
'e'):
636 elif tag.startswith(
's'):
638 elif tag.startswith(
'r'):
640 elif tag.startswith(
'p'):
642 elif tag.startswith(
'a'):
644 elif tag.startswith(
't'):
646 elif tag.startswith(
'b'):
651 fileParts = filePart.split(
'.')
652 result[
'jediTaskID'] = fileParts[1]
653 result[
'fileNumber'] = fileParts[2]
654 result[
'version'] = fileParts[-1]
660 with path.open(
'r')
as inputText:
661 for line
in inputText.readlines():
663 if line.strip().startswith(
"#")
or not line.strip():
665 files += line.split(
",")
667 files = [file.strip()
for file
in files]
671 if any((path.parent / file).
exists()
or (path.parent / f
"{file}.txt").
exists()
for file
in files):
675 file_path = path.parent / file
676 if not file_path.exists():
677 file_path = path.parent / f
"{file}.txt"
678 if not file_path.exists():
679 logCPGridRun.error(f
"File {file} or {file}.txt does not exist in the input list directory.")
680 raise FileNotFoundError(f
"File {file} or {file}.txt does not exist in the input list directory.")
681 files_current, names_current = CPGridRun._parseInputFileList(file_path, bulk_submission=
True)
682 files_bulk.extend(files_current)
683 names_bulk.extend(names_current)
684 return files_bulk, names_bulk
686 return [
','.join(files)], [path.stem.replace(
"+",
"")]
688 return files, [
None] * len(files)
692 logCPGridRun.error(
"Errors were collected during the script execution:")
695 logCPGridRun.error(f
"{key}: {value}")
696 logCPGridRun.error(
"Please fix the errors and try again.")
701 if self.
args.checkInputDS:
705 if self.
args.agreeAll:
706 logCPGridRun.info(
"You have agreed to all the submission details. Jobs will be submitted without confirmation.")
709 answer = input(
"Please confirm ALL the submission details are correct before submitting [y/n]: ")
710 if answer.lower() ==
'y':
712 elif answer.lower() ==
'n':
713 logCPGridRun.info(
"Feel free to report any unexpected behavior to the CPAlgorithms team!")
715 logCPGridRun.error(
"Invalid input. Please enter 'y' or 'n'. Jobs are not submitted.")
717if __name__ ==
'__main__':
719 cpgrid.configureSubmission()
720 cpgrid.printInputDetails()
721 cpgrid.checkExternalTools()
722 cpgrid.printDelayedErrorCollection()
723 cpgrid.askSubmission()
void print(char *figname, TCanvas *c1)
dict _createPrunArgsDict(self)
outputDSFormatter(self, name, label)
rucioCustomNameParser(filename)
_checkYamlExists(self, runscriptArgs)
_parseGridArguments(self)
bool checkInputInPyami(self)
bool _analyzeAmiResults(self, results, datasetPtag)
isAtlasProductionFormat(name)
tuple[list[str], list[str]] _parseInputFileList(Path path, bool bulk_submission=False)
_prepareAmiQueryFromInputList(self)
configureSubmissionSingleSample(self, input, name)
_filesChangedOrTarballNotCreated(self)
tuple[list[str], list[str]] _inputNames
_hasCompressedTarball(self)
atlasProductionNameParser(filename)
_outputDSFormatter(self, name)
_checkPrunArgs(self, argDict)
configureSubmission(self)
printDelayedErrorCollection(self)
_customOutputDSFormatter(self, name, label)
dict _unknownArgsDict(self)
bool exists(const std::string &filename)
does a file exist
std::vector< std::string > split(const std::string &s, const std::string &t=":")