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));
511 using namespace msgEventLoop;
517 const std::string jobELGDir = data.submitDir +
"/elg";
518 const std::string runShFile = jobELGDir +
"/runjob.sh";
520 const std::string mergeShFile = jobELGDir +
"/elg_merge";
526 const std::string jobDefFile = jobELGDir +
"/jobdef.root";
527 gSystem->Exec(Form(
"mkdir -p %s", jobELGDir.c_str()));
528 gSystem->Exec(Form(
"cp %s %s", runShOrig.c_str(), runShFile.c_str()));
529 gSystem->Exec(Form(
"chmod +x %s", runShFile.c_str()));
530 gSystem->Exec(Form(
"cp %s %s", mergeShOrig.c_str(), mergeShFile.c_str()));
531 gSystem->Exec(Form(
"chmod +x %s", mergeShFile.c_str()));
536 if (listToShipToGrid.size()){
538 "Creating symbolic links for additional files or directories to be sent to grid.\n"
539 "For root or heavy files you should also add their name (not the full path) to EL::Job::optUserFiles.\n"
540 "Otherwise prun ignores those files."
543 std::vector<std::string> vect_filesOrDirToShip;
544 for (
auto&& part : std::views::split(listToShipToGrid,
',')) vect_filesOrDirToShip.emplace_back(part.begin(), part.end());
546 for (
const std::string & fileOrDirToShip: vect_filesOrDirToShip){
547 ANA_MSG_INFO ((
"Creating symbolic link for: " +fileOrDirToShip).c_str());
557 meta.fetchDefaults(data.options);
560 std::string outputSampleName =
meta.castString(
"nc_outputSampleName");
561 if (outputSampleName.empty()) {
562 outputSampleName =
"user.%nickname%.%in:name%";
565 meta.setString(
"nc_inDS",
meta.castString(
"nc_grid", (*s)->name()));
566 meta.setString(
"nc_writeInputToTxt",
"IN:input.txt");
567 meta.setString(
"nc_match",
meta.castString(
"nc_grid_filter"));
568 const std::string execstr =
"runjob.sh " + (*s)->name();
569 meta.setString(
"nc_exec", execstr);
570 meta.setString(
"nc_framework",
"EventLoopGrid");
576 out != data.job->outputEnd(); ++out) {
578 shOut.
save(data.submitDir +
"/output-" + out->label());
581 shHist.
save(data.submitDir +
"/output-hist");
583 TmpCd keepDir(jobELGDir);
587 sh.save(data.submitDir +
"/input");
588 data.submitted =
true;
601 return ::StatusCode::SUCCESS;
609 TmpCd tmpDir(data.submitDir);
615 const size_t nRunThreads =
options()->castDouble(
"nc_run_threads", 0);
616 const size_t nDlThreads =
options()->castDouble(
"nc_download_threads", 0);
624 std::cout << std::endl;
632 std::cout << (*s)->name() <<
"\t";
636 case JobState::DOWNLOAD:
637 case JobState::MERGE:
638 std::cout << JobState::name[state] <<
"\t";
640 case JobState::FINISHED:
641 std::cout <<
"\033[1;32m" << JobState::name[state] <<
"\033[0m\t";
643 case JobState::FAILED:
644 std::cout <<
"\033[1;31m" << JobState::name[state] <<
"\033[0m\t";
647 std::cout <<
details << std::endl;
649 allDone &= (state == JobState::FINISHED || state == JobState::FAILED);
652 std::cout << std::endl;
654 data.retrieved =
true;
655 data.completed = allDone;
656 return ::StatusCode::SUCCESS;
662 TmpCd tmpDir(location);
672 std::cout << (*s)->name() <<
"\t" << JobState::name[state]
673 <<
"\t" <<
details << std::endl;
678 const std::string& task,
679 const std::string& state)
684 TmpCd tmpDir(location);
688 if (not
sh.get(task)) {
689 std::cout <<
"Unknown task: " << task << std::endl;
690 std::cout <<
"Choose one of: " << std::endl;
694 JobState::parse(state);
695 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