51 static const unsigned int NSTATES = 6;
52 enum Enum {
INIT=0, RUN=1, DOWNLOAD=2, MERGE=3, FINISHED=4, FAILED=5 };
53 static const char* name[NSTATES] =
54 {
"INIT",
"RUNNING",
"DOWNLOAD",
"MERGE",
"FINISHED",
"FAILED" };
55 Enum
parse(
const std::string& what)
57 for (
unsigned int i = 0; i != NSTATES; ++i) {
58 if (what == name[i]) {
return static_cast<Enum
>(i); }
61 throw std::runtime_error(
"PrunDriver.cxx: Failed to parse job state string");
70 enum Enum { DONE=0, PENDING=1, FAIL=2 };
73 struct TransitionRule {
74 JobState::Enum fromState;
76 JobState::Enum toState;
77 TransitionRule(JobState::Enum fromState,
79 JobState::Enum toState)
80 : fromState(fromState)
88 const std::string origDir;
89 TmpCd(
const std::string & dir)
90 : origDir(gSystem->pwd())
92 gSystem->cd(dir.c_str());
96 gSystem->cd(origDir.c_str());
104 static const std::string defaultState = JobState::name[JobState::INIT];
106 return JobState::parse(
label);
109static JobState::Enum
nextState(JobState::Enum state, Status::Enum status)
113 static const TransitionRule TABLE[] =
115 TransitionRule(JobState::INIT, Status::DONE, JobState::RUN),
116 TransitionRule(JobState::INIT, Status::PENDING, JobState::INIT),
117 TransitionRule(JobState::INIT, Status::FAIL, JobState::FAILED),
118 TransitionRule(JobState::RUN, Status::DONE, JobState::DOWNLOAD),
119 TransitionRule(JobState::RUN, Status::PENDING, JobState::RUN),
120 TransitionRule(JobState::RUN, Status::FAIL, JobState::FAILED),
121 TransitionRule(JobState::DOWNLOAD, Status::DONE, JobState::MERGE),
122 TransitionRule(JobState::DOWNLOAD, Status::PENDING, JobState::DOWNLOAD),
123 TransitionRule(JobState::DOWNLOAD, Status::FAIL, JobState::FAILED),
124 TransitionRule(JobState::MERGE, Status::DONE, JobState::FINISHED),
125 TransitionRule(JobState::MERGE, Status::PENDING, JobState::MERGE),
126 TransitionRule(JobState::MERGE, Status::FAIL, JobState::DOWNLOAD)
128 static const unsigned int TABLE_SIZE =
sizeof(TABLE) /
sizeof(TABLE[0]);
129 for (
unsigned int i = 0; i != TABLE_SIZE; ++i) {
130 if (TABLE[i].fromState == state && TABLE[i].status == status) {
131 return TABLE[i].toState;
135 throw std::logic_error(
"PrunDriver.cxx: Missing state transition rule");
143 o.
setString(
"nc_cmtConfig", gSystem->ExpandPathName(
"$AnalysisBase_PLATFORM"));
144 o.
setString(
"nc_useAthenaPackages",
"true");
145 const std::string mergestr =
"elg_merge jobdef.root %OUT %IN";
151 const std::string& location)
158 gSystem->Exec(Form(
"mkdir -p %s", location.c_str()));
160 std::vector<std::string> datasets;
163 if (entry.type ==
"CONTAINER" || entry.type ==
"DIDType.CONTAINER")
164 datasets.push_back (entry.name);
168 for (
const auto& result : downloadResult)
170 if (result.notDownloaded != 0)
182 using namespace EL::msgEventLoop;
184 ANA_MSG_INFO(
"Submitting " << sample->name() <<
"..." );
186 static bool loaded =
false;
192 TPython::LoadMacro(path.c_str());
196 TPython::Bind(
dynamic_cast<TObject*
>(sample),
"ELG_SAMPLE");
197#if ROOT_VERSION_CODE >= ROOT_VERSION(6,33,01)
199 TPython::Exec(
"_anyresult = ROOT.std.make_any['int'](ELG_prun(ELG_SAMPLE))", &result);
200 int ret = std::any_cast<int>(result);
202 int ret = TPython::Eval(
"ELG_prun(ELG_SAMPLE)");
204 TPython::Bind(0,
"ELG_SAMPLE");
211 if (isFirstSample && ret == 1){
213 throw std::runtime_error(
"PrunDriver.cxx: aborting due to tarball creation issue");
217 sample->meta()->setString(
"nc_ELG_state_details",
218 "problem submitting");
222 sample->meta()->setDouble(
"nc_jediTaskID", ret);
235 static bool loaded =
false;
238 TPython::LoadMacro(path.c_str());
242 TPython::Bind(
dynamic_cast<TObject*
>(sample),
"ELG_SAMPLE");
243#if ROOT_VERSION_CODE >= ROOT_VERSION(6,33,01)
245 TPython::Exec(
"_anyresult = ROOT.std.make_any['int'](ELG_jediState(ELG_SAMPLE))", &result);
246 int ret = std::any_cast<int>(result);
248 int ret = TPython::Eval(
"ELG_jediState(ELG_SAMPLE)");
250 TPython::Bind(0,
"ELG_SAMPLE");
252 if (ret == Status::DONE)
return Status::DONE;
253 if (ret == Status::FAIL)
return Status::FAIL;
256 if (ret != 90) { sample->meta()->setString(
"nc_ELG_state_details",
"task status other than done/finished/failed/running"); }
259 sample->meta()->setString(
"nc_ELG_state_details",
260 "problem checking jedi task status");
263 return Status::PENDING;
271 static std::mutex
mutex;
273 std::cout <<
"Downloading output from: "
274 << sample->name() <<
"..." << std::endl;
279 if (container[container.size()-1] ==
'/') {
280 container.resize(container.size() - 1);
282 container +=
"_hist/";
286 if (not downloadOk) {
287 std::cerr <<
"Failed to download one or more files" << std::endl;
288 sample->meta()->setString(
"nc_ELG_state_details",
289 "error, check log for details");
290 return Status::PENDING;
302 if (container[container.size()-1] ==
'/') {
303 container.resize(container.size() - 1);
305 container +=
"_hist/";
306 const std::string dir =
"elg/download/" + container;
308 const std::string fileName =
"hist-output.root";
310 const std::string target = Form(
"hist-%s.root", sample->name().c_str());
312 const std::string findCmd(Form(
"find %s -name \"*.%s*\" | tr '\n' ' '",
313 dir.c_str(), fileName.c_str()));
314 std::istringstream input(gSystem->GetFromPipe(findCmd.c_str()).Data());
315 std::vector<std::string>
files((std::istream_iterator<std::string>(input)),
316 std::istream_iterator<std::string>());
321 if (not
files.size()) {
322 std::cerr <<
"Found no input files for merging! "
323 <<
"Requeueing sample for download..." << std::endl;
324 sample->meta()->setString(
"nc_ELG_state_details",
"retry, files were lost");
331 sample->meta()->setString(
"nc_ELG_state_details",
332 "error, check log for details");
333 gSystem->Exec(Form(
"rm -f %s", target.c_str()));
334 return Status::PENDING;
337 for (
size_t i = 0; i !=
files.size(); ++i) {
338 gSystem->Exec(Form(
"rm %s",
files[i].c_str()));
340 gSystem->Exec(Form(
"rmdir %s/*", dir.c_str()));
341 gSystem->Exec(Form(
"rmdir %s", dir.c_str()));
352 sample->meta()->setString(
"nc_ELG_state_details",
"");
354 Status::Enum status = Status::PENDING;
357 status =
submit(sample, isFirstSample);
362 case JobState::DOWNLOAD:
365 case JobState::MERGE:
366 status =
merge(sample);
368 case JobState::FINISHED:
369 case JobState::FAILED:
374 sample->meta()->setString(
"nc_ELG_state", JobState::name[state]);
378 const size_t nThreads)
384 bool isFirstSample =
true;
387 workList.push_back([s, isFirstSample]()->
void{
processTask(*s, isFirstSample); });
389 isFirstSample =
false;
396 const std::string & pattern)
398 const std::string sampleName = sampleMeta.
castString(
"sample_name");
400 using namespace EL::msgEventLoop;
402 static const std::string nickname =
403 gSystem->GetFromPipe(Form(
"python -c \"%s\" 2>/dev/null",
404 "from pandatools import PsubUtils;"
405 "print(PsubUtils.getNickname());")).Data();
407 TString out = pattern.c_str();
410 if (nickname.length()>20){
411 ANA_MSG_WARNING(
"No proxy available - cannot use nickname yet. Will try a late replacement.");
413 out.ReplaceAll(
"%nickname%", nickname);
416 out.ReplaceAll(
"%in:name%", sampleName);
418 std::stringstream
ss(sampleName);
421 while(std::getline(
ss, item,
'.')) {
422 std::stringstream sskey;
423 sskey <<
"%in:name[" << ++field <<
"]%";
424 out.ReplaceAll(sskey.str(), item);
426 while (out.Index(
"%in:") != -1) {
427 int i1 = out.Index(
"%in:");
428 int i2 = out.Index(
"%", i1+1);
429 TString metaName = out(i1+4, i2-i1-4);
430 out.ReplaceAll(
"%in:"+metaName+
"%",
431 sampleMeta.
castString(std::string(metaName.Data())));
433 out.ReplaceAll(
"/",
"");
441 end = job.outputEnd(); out != end; ++out) {
442 outputs.Add(out->Clone());
444 std::string out =
"hist:hist-output.root";
447 while ((obj = itr())) {
449 const std::string name = os->label() +
".root";
450 const std::string ds =
452 out +=
"," + (ds.empty() ? name : ds +
":" + name);
462 TFile
file(fileName.c_str(),
"RECREATE");
465 outputs.Add(o->Clone());
466 file.WriteTObject(&job.jobConfig(),
"jobConfig",
"SingleKey");
467 file.WriteTObject(&outputs,
"outputs",
"SingleKey");
468 bool haveDefault =
false;
471 file.WriteObject(&
meta,
meta.castString(
"sample_name").c_str());
474 file.WriteObject (&
meta,
"defaultMetaObject");
485 const std::string outputFile =
"*" +
outputLabel +
".root*";
486 const std::string outDSSuffix =
'_' +
outputLabel +
".root/";
488 auto outSample = std::make_unique<SH::SampleGrid>((*s)->name());
490 outSample->meta()->setString(
"nc_grid", outputDS);
491 outSample->meta()->setString(
"nc_grid_filter", outputFile);
492 out.add(std::move(outSample));
509 using namespace msgEventLoop;
515 const std::string jobELGDir = data.submitDir +
"/elg";
516 const std::string runShFile = jobELGDir +
"/runjob.sh";
518 const std::string mergeShFile = jobELGDir +
"/elg_merge";
524 const std::string jobDefFile = jobELGDir +
"/jobdef.root";
525 gSystem->Exec(Form(
"mkdir -p %s", jobELGDir.c_str()));
526 gSystem->Exec(Form(
"cp %s %s", runShOrig.c_str(), runShFile.c_str()));
527 gSystem->Exec(Form(
"chmod +x %s", runShFile.c_str()));
528 gSystem->Exec(Form(
"cp %s %s", mergeShOrig.c_str(), mergeShFile.c_str()));
529 gSystem->Exec(Form(
"chmod +x %s", mergeShFile.c_str()));
534 if (listToShipToGrid.size()){
536 "Creating symbolic links for additional files or directories to be sent to grid.\n"
537 "For root or heavy files you should also add their name (not the full path) to EL::Job::optUserFiles.\n"
538 "Otherwise prun ignores those files."
541 std::vector<std::string> vect_filesOrDirToShip;
542 for (
auto&& part : std::views::split(listToShipToGrid,
',')) vect_filesOrDirToShip.emplace_back(part.begin(), part.end());
544 for (
const std::string & fileOrDirToShip: vect_filesOrDirToShip){
545 ANA_MSG_INFO ((
"Creating symbolic link for: " +fileOrDirToShip).c_str());
555 meta.fetchDefaults(data.options);
558 std::string outputSampleName =
meta.castString(
"nc_outputSampleName");
559 if (outputSampleName.empty()) {
560 outputSampleName =
"user.%nickname%.%in:name%";
563 meta.setString(
"nc_inDS",
meta.castString(
"nc_grid", (*s)->name()));
564 meta.setString(
"nc_writeInputToTxt",
"IN:input.txt");
565 meta.setString(
"nc_match",
meta.castString(
"nc_grid_filter"));
566 const std::string execstr =
"runjob.sh " + (*s)->name();
567 meta.setString(
"nc_exec", execstr);
568 meta.setString(
"nc_framework",
"EventLoopGrid");
574 out != data.job->outputEnd(); ++out) {
576 shOut.
save(data.submitDir +
"/output-" + out->label());
579 shHist.
save(data.submitDir +
"/output-hist");
581 TmpCd keepDir(jobELGDir);
585 sh.save(data.submitDir +
"/input");
586 data.submitted =
true;
599 return ::StatusCode::SUCCESS;
607 TmpCd tmpDir(data.submitDir);
613 const size_t nRunThreads =
options()->castDouble(
"nc_run_threads", 0);
614 const size_t nDlThreads =
options()->castDouble(
"nc_download_threads", 0);
622 std::cout << std::endl;
630 std::cout << (*s)->name() <<
"\t";
634 case JobState::DOWNLOAD:
635 case JobState::MERGE:
636 std::cout << JobState::name[state] <<
"\t";
638 case JobState::FINISHED:
639 std::cout <<
"\033[1;32m" << JobState::name[state] <<
"\033[0m\t";
641 case JobState::FAILED:
642 std::cout <<
"\033[1;31m" << JobState::name[state] <<
"\033[0m\t";
645 std::cout <<
details << std::endl;
647 allDone &= (state == JobState::FINISHED || state == JobState::FAILED);
650 std::cout << std::endl;
652 data.retrieved =
true;
653 data.completed = allDone;
654 return ::StatusCode::SUCCESS;
660 TmpCd tmpDir(location);
670 std::cout << (*s)->name() <<
"\t" << JobState::name[state]
671 <<
"\t" <<
details << std::endl;
676 const std::string& task,
677 const std::string& state)
682 TmpCd tmpDir(location);
686 if (not
sh.get(task)) {
687 std::cout <<
"Unknown task: " << task << std::endl;
688 std::cout <<
"Choose one of: " << std::endl;
692 JobState::parse(state);
693 sh.get(task)->meta()->setString(
"nc_ELG_state", state);
#define RCU_NEW_INVARIANT(x)
#define RCU_READ_INVARIANT(x)
virtual void lock()=0
Interface to allow an object to lock itself when made const in SG.
const std::string outputLabel
std::string PathResolverFindCalibFile(const std::string &logical_file_name)
static SH::MetaObject defaultOpts()
static void processAllInState(const SH::SampleHandler &sh, JobState::Enum state, const size_t nThreads)
ClassImp(EL::PrunDriver) namespace
std::string outputFileNames(const EL::Job &job)
static Status::Enum submit(SH::Sample *const sample, const bool isFirstSample)
static bool downloadContainer(const std::string &name, const std::string &location)
static JobState::Enum nextState(JobState::Enum state, Status::Enum status)
static std::string formatOutputName(const SH::MetaObject &sampleMeta, const std::string &pattern)
static Status::Enum checkPandaTask(SH::Sample *const sample)
static void saveJobDef(const std::string &fileName, const EL::Job &job, const SH::SampleHandler sh)
static Status::Enum download(SH::Sample *const sample)
static JobState::Enum sampleState(SH::Sample *sample)
static SH::SampleHandler outputSH(const SH::SampleHandler &in, const std::string &outputLabel)
static void processTask(SH::Sample *const sample, const bool isFirstSample)
SH::MetaObject * options()
the list of options to jobs with this driver
virtual::StatusCode doManagerStep(Detail::ManagerData &data) const
const OutputStream * outputIter
static const std::string optGridPrunShipAdditionalFilesOrDirs
Enables to ship additional files to the tarbal sent to the grid Should be a list of comma separated p...
static const std::string optContainerSuffix
a Driver to submit jobs via prun
static void status(const std::string &location)
::StatusCode doRetrieve(Detail::ManagerData &data) const
static void setState(const std::string &location, const std::string &task, const std::string &state)
void testInvariant() const
A class that manages a list of Sample objects.
void save(const std::string &directory) const
save the list of samples to the given directory
iterator begin() const
the begin iterator to use
boost::transform_iterator< SamplePtrToRawSample, std::vector< std::shared_ptr< Sample > >::const_iterator > iterator
the iterator to use
iterator end() const
the end iterator to use
a base class that manages a set of files belonging to a particular data set and the associated meta-d...
const std::string process
std::vector< std::string > files
file names and file pointers
std::string label(const std::string &format, int i)
@ doRetrieve
call the actual doRetrieve method
@ submitJob
do the actual job submission
::StatusCode StatusCode
StatusCode definition for legacy code.
void exec(const std::string &cmd)
effects: execute the given command guarantee: strong failures: out of memory II failures: system fail...
void hadd(const std::string &output_file, const std::vector< std::string > &input_files, unsigned max_files)
effects: perform the hadd functionality guarantee: basic failures: out of memory III failures: i/o er...
std::vector< RucioDownloadResult > rucioDownloadList(const std::string &location, const std::vector< std::string > &datasets)
run rucio-download with multiple datasets
std::vector< RucioListDidsEntry > rucioListDids(const std::string &dataset)
run rucio-list-dids for the given dataset
DataModel_detail::iterator< DVL > unique(typename DataModel_detail::iterator< DVL > beg, typename DataModel_detail::iterator< DVL > end)
Specialization of unique for DataVector/List.
void sort(typename DataModel_detail::iterator< DVL > beg, typename DataModel_detail::iterator< DVL > end)
Specialization of sort for DataVector/List.
an internal data structure for passing data between different manager objects anbd step