ATLAS Offline Software
Loading...
Searching...
No Matches
PrunDriver.cxx File Reference
#include <EventLoopGrid/PrunDriver.h>
#include <EventLoop/Algorithm.h>
#include <EventLoop/ManagerData.h>
#include <EventLoop/ManagerStep.h>
#include <EventLoop/Job.h>
#include <EventLoop/MessageCheck.h>
#include <EventLoop/OutputStream.h>
#include <PathResolver/PathResolver.h>
#include <RootCoreUtils/Assert.h>
#include <RootCoreUtils/hadd.h>
#include <RootCoreUtils/ShellExec.h>
#include <SampleHandler/MetaObject.h>
#include <SampleHandler/Sample.h>
#include <SampleHandler/SampleGrid.h>
#include <SampleHandler/SampleHandler.h>
#include <SampleHandler/GridTools.h>
#include <TList.h>
#include <TPython.h>
#include <TROOT.h>
#include <TFile.h>
#include <TSystem.h>
#include <algorithm>
#include <cstdlib>
#include <iostream>
#include <sstream>
#include <string>
#include <vector>
#include <stdexcept>
#include <ranges>
#include "pool.h"
#include <mutex>

Go to the source code of this file.

Functions

 ClassImp (EL::PrunDriver) namespace
static JobState::Enum sampleState (SH::Sample *sample)
static JobState::Enum nextState (JobState::Enum state, Status::Enum status)
static SH::MetaObject defaultOpts ()
static bool downloadContainer (const std::string &name, const std::string &location)
static Status::Enum submit (SH::Sample *const sample, const bool isFirstSample)
static Status::Enum checkPandaTask (SH::Sample *const sample)
static Status::Enum download (SH::Sample *const sample)
static Status::Enum merge (SH::Sample *const sample)
static void processTask (SH::Sample *const sample, const bool isFirstSample)
static void processAllInState (const SH::SampleHandler &sh, JobState::Enum state, const size_t nThreads)
static std::string formatOutputName (const SH::MetaObject &sampleMeta, const std::string &pattern)
std::string outputFileNames (const EL::Job &job)
static void saveJobDef (const std::string &fileName, const EL::Job &job, const SH::SampleHandler sh)
static SH::SampleHandler outputSH (const SH::SampleHandler &in, const std::string &outputLabel)

Function Documentation

◆ checkPandaTask()

Status::Enum checkPandaTask ( SH::Sample *const sample)
static

Definition at line 230 of file PrunDriver.cxx.

231{
232 RCU_REQUIRE(sample);
233 RCU_REQUIRE(static_cast<int>(sample->meta()->castDouble("nc_jediTaskID",0, SH::MetaObject::CAST_NOCAST_DEFAULT)) > 100);
234
235 static bool loaded = false;
236 if (not loaded) {
237 std::string path = PathResolverFindCalibFile("EventLoopGrid/ELG_jediState.py");
238 TPython::LoadMacro(path.c_str());
239 loaded = true;
240 }
241
242 TPython::Bind(dynamic_cast<TObject*>(sample), "ELG_SAMPLE");
243#if ROOT_VERSION_CODE >= ROOT_VERSION(6,33,01)
244 std::any result;
245 TPython::Exec("_anyresult = ROOT.std.make_any['int'](ELG_jediState(ELG_SAMPLE))", &result);
246 int ret = std::any_cast<int>(result);
247#else
248 int ret = TPython::Eval("ELG_jediState(ELG_SAMPLE)");
249#endif
250 TPython::Bind(0, "ELG_SAMPLE");
251
252 if (ret == Status::DONE) return Status::DONE;
253 if (ret == Status::FAIL) return Status::FAIL;
254
255 // Value 90 corresponds to `running` state of the job
256 if (ret != 90) { sample->meta()->setString("nc_ELG_state_details", "task status other than done/finished/failed/running"); }
257 // Value 99 is returned if there is error in the script (import, missing ID)
258 if (ret == 99) {
259 sample->meta()->setString("nc_ELG_state_details",
260 "problem checking jedi task status");
261 }
262
263 return Status::PENDING;
264}
#define RCU_REQUIRE(x)
Definition Assert.h:196
std::string PathResolverFindCalibFile(const std::string &logical_file_name)
@ CAST_NOCAST_DEFAULT
cast and return the default value if the input has the wrong type
Definition MetaObject.h:70
path
python interpreter configuration --------------------------------------—
Definition athena.py:130

◆ ClassImp()

ClassImp ( EL::PrunDriver )
Author
Alexander Madsen
Nils Krumnack

Definition at line 47 of file PrunDriver.cxx.

49 {
50 namespace JobState {
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)
56 {
57 for (unsigned int i = 0; i != NSTATES; ++i) {
58 if (what == name[i]) { return static_cast<Enum>(i); }
59 }
60 RCU_ASSERT0("Failed to parse job state string");
61 throw std::runtime_error("PrunDriver.cxx: Failed to parse job state string"); //compiler dummy
62 }
63 }
64
65 // When changing the values in the enum make sure
66 // corresponding values in `data/ELG_jediState.py` script
67 // are changed accordingly
68 namespace Status {
69 //static const int NSTATES = 3;
70 enum Enum { DONE=0, PENDING=1, FAIL=2 };
71 }
72
73 struct TransitionRule {
74 JobState::Enum fromState;
75 Status::Enum status;
76 JobState::Enum toState;
77 TransitionRule(JobState::Enum fromState,
78 Status::Enum status,
79 JobState::Enum toState)
80 : fromState(fromState)
81 , status(status)
82 , toState(toState)
83 {
84 }
85 };
86
87 struct TmpCd {
88 const std::string origDir;
89 TmpCd(const std::string & dir)
90 : origDir(gSystem->pwd())
91 {
92 gSystem->cd(dir.c_str());
93 }
94 ~TmpCd()
95 {
96 gSystem->cd(origDir.c_str());
97 }
98 };
99}
#define RCU_ASSERT0(y)
Definition Assert.h:214
#define INIT(__TYPE)
std::map< std::string, std::string, std::less<> > parse(const std::string &list)
Status
Athena specific StatusCode values.
status
Definition merge.py:16

◆ defaultOpts()

SH::MetaObject defaultOpts ( )
static

Definition at line 138 of file PrunDriver.cxx.

139{
141 o.setString("nc_nGBPerJob", "MAX");
142 o.setString("nc_mergeOutput", "true");
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";
146 o.setString("nc_mergeScript", mergestr);
147 return o;
148}
A class that manages meta-data to be associated with an object.
Definition MetaObject.h:48
void setString(const std::string &name, const std::string &value)
set the meta-data string with the given name

◆ download()

Status::Enum download ( SH::Sample *const sample)
static

Definition at line 266 of file PrunDriver.cxx.

267{
268 RCU_REQUIRE(sample);
269
270 {
271 static std::mutex mutex;
272 std::lock_guard<std::mutex> lock(mutex);
273 std::cout << "Downloading output from: "
274 << sample->name() << "..." << std::endl;
275 }
276
277 std::string container = sample->meta()->castString("nc_outDS", "", SH::MetaObject::CAST_NOCAST_DEFAULT);
278 RCU_ASSERT(not container.empty());
279 if (container[container.size()-1] == '/') {
280 container.resize(container.size() - 1);
281 }
282 container += "_hist/";
283
284 bool downloadOk = downloadContainer(container, "elg/download/" + container);
285
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;
291 }
292
293 return Status::DONE;
294}
#define RCU_ASSERT(x)
Definition Assert.h:210
virtual void lock()=0
Interface to allow an object to lock itself when made const in SG.
static bool downloadContainer(const std::string &name, const std::string &location)
STL class.
const SG::AuxVectorData * container() const
Return the container holding this element.

◆ downloadContainer()

bool downloadContainer ( const std::string & name,
const std::string & location )
static

Definition at line 150 of file PrunDriver.cxx.

152{
153 RCU_ASSERT(not name.empty());
154 RCU_ASSERT(name[name.size()-1] == '/');
155 RCU_ASSERT(not location.empty());
156
157 try {
158 gSystem->Exec(Form("mkdir -p %s", location.c_str()));
159
160 std::vector<std::string> datasets;
161 for (auto& entry : SH::rucioListDids (name))
162 {
163 if (entry.type == "CONTAINER" || entry.type == "DIDType.CONTAINER")
164 datasets.push_back (entry.name);
165 }
166
167 auto downloadResult = SH::rucioDownloadList (location, datasets);
168 for (const auto& result : downloadResult)
169 {
170 if (result.notDownloaded != 0)
171 return false;
172 }
173 } catch (...) {
174 return false;
175 }
176 return true;
177}
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

◆ formatOutputName()

std::string formatOutputName ( const SH::MetaObject & sampleMeta,
const std::string & pattern )
static

Definition at line 395 of file PrunDriver.cxx.

397{
398 const std::string sampleName = sampleMeta.castString("sample_name");
399 RCU_REQUIRE(not pattern.empty());
400 using namespace EL::msgEventLoop;
401
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();
406
407 TString out = pattern.c_str();
408
409 // Handle case of no proxy; will create a proxy later in the submission
410 if (nickname.length()>20){
411 ANA_MSG_WARNING( "No proxy available - cannot use nickname yet. Will try a late replacement.");
412 } else {
413 out.ReplaceAll("%nickname%", nickname);
414 }
415
416 out.ReplaceAll("%in:name%", sampleName);
417
418 std::stringstream ss(sampleName);
419 std::string item;
420 int field = 0;
421 while(std::getline(ss, item, '.')) {
422 std::stringstream sskey;
423 sskey << "%in:name[" << ++field << "]%";
424 out.ReplaceAll(sskey.str(), item);
425 }
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())));
432 }
433 out.ReplaceAll("/", "");
434 return out.Data();
435}
#define ANA_MSG_WARNING(xmsg,...)
Macro printing warning messages.
static Double_t ss
std::string castString(const std::string &name, const std::string &def_val="", CastMode mode=CAST_ERROR_THROW) const
the meta-data string with the given name

◆ merge()

Status::Enum merge ( SH::Sample *const sample)
static

Definition at line 296 of file PrunDriver.cxx.

297{
298 RCU_REQUIRE(sample);
299
300 std::string container = sample->meta()->castString("nc_outDS", "", SH::MetaObject::CAST_NOCAST_DEFAULT);
301 RCU_ASSERT(not container.empty());
302 if (container[container.size()-1] == '/') {
303 container.resize(container.size() - 1);
304 }
305 container += "_hist/";
306 const std::string dir = "elg/download/" + container;
307
308 const std::string fileName = "hist-output.root";
309
310 const std::string target = Form("hist-%s.root", sample->name().c_str());
311
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>());
317
318 std::sort(files.begin(), files.end());
319 RCU_ASSERT(std::unique(files.begin(), files.end()) == files.end());
320
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");
325 return Status::FAIL;
326 }
327
328 try {
329 RCU::hadd(target.c_str(), files);
330 } catch (...) {
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;
335 }
336
337 for (size_t i = 0; i != files.size(); ++i) {
338 gSystem->Exec(Form("rm %s", files[i].c_str()));
339 }
340 gSystem->Exec(Form("rmdir %s/*", dir.c_str()));
341 gSystem->Exec(Form("rmdir %s", dir.c_str()));
342
343 return Status::DONE;
344}
std::vector< std::string > files
file names and file pointers
Definition hcg.cxx:52
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...
Definition hadd.cxx:29
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.

◆ nextState()

JobState::Enum nextState ( JobState::Enum state,
Status::Enum status )
static

Definition at line 109 of file PrunDriver.cxx.

110{
111 RCU_REQUIRE(state != JobState::FINISHED);
112 RCU_REQUIRE(state != JobState::FAILED);
113 static const TransitionRule TABLE[] =
114 {
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)
127 };
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;
132 }
133 }
134 RCU_ASSERT0("Missing state transition rule");
135 throw std::logic_error("PrunDriver.cxx: Missing state transition rule");
136}

◆ outputFileNames()

std::string outputFileNames ( const EL::Job & job)

Definition at line 437 of file PrunDriver.cxx.

438{
439 TList outputs;
440 for (EL::Job::outputIter out = job.outputBegin(),
441 end = job.outputEnd(); out != end; ++out) {
442 outputs.Add(out->Clone());
443 }
444 std::string out = "hist:hist-output.root";
445 TIter itr(&outputs);
446 TObject *obj = 0;
447 while ((obj = itr())) {
448 EL::OutputStream *os = dynamic_cast<EL::OutputStream*>(obj);
449 const std::string name = os->label() + ".root";
450 const std::string ds =
451 os->options()->castString(EL::OutputStream::optContainerSuffix);
452 out += "," + (ds.empty() ? name : ds + ":" + name);
453 }
454 return out;
455}
const OutputStream * outputIter
Definition Job.h:139
static const std::string optContainerSuffix

◆ outputSH()

SH::SampleHandler outputSH ( const SH::SampleHandler & in,
const std::string & outputLabel )
static

Definition at line 481 of file PrunDriver.cxx.

483{
485 const std::string outputFile = "*" + outputLabel + ".root*";
486 const std::string outDSSuffix = '_' + outputLabel + ".root/";
487 for (SH::SampleHandler::iterator s = in.begin(); s != in.end(); ++s) {
488 auto outSample = std::make_unique<SH::SampleGrid>((*s)->name());
489 const std::string outputDS = (*s)->meta()->castString("nc_outDS", "", SH::MetaObject::CAST_NOCAST_DEFAULT) + outDSSuffix;
490 outSample->meta()->setString("nc_grid", outputDS);
491 outSample->meta()->setString("nc_grid_filter", outputFile);
492 out.add(std::move(outSample));
493 }
494 out.fetch(in);
495 return out;
496}
const std::string outputLabel
A class that manages a list of Sample objects.
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

◆ processAllInState()

void processAllInState ( const SH::SampleHandler & sh,
JobState::Enum state,
const size_t nThreads )
static

Definition at line 377 of file PrunDriver.cxx.

379{
380 RCU_REQUIRE(sh.size());
381
382 WorkList workList;
383
384 bool isFirstSample = true;
385 for (SH::SampleHandler::iterator s = sh.begin(); s != sh.end(); ++s) {
386 if (sampleState(*s) == state) {
387 workList.push_back([s, isFirstSample]()->void{ processTask(*s, isFirstSample); });
388 // Change boolean to false as already processed one sample
389 isFirstSample = false;
390 }
391 }
392 process(workList, nThreads);
393}
std::vector< WorkUnit > WorkList
static JobState::Enum sampleState(SH::Sample *sample)
static void processTask(SH::Sample *const sample, const bool isFirstSample)
const std::string process

◆ processTask()

void processTask ( SH::Sample *const sample,
const bool isFirstSample )
static

Definition at line 346 of file PrunDriver.cxx.

347{
348 RCU_REQUIRE(sample);
349
350 JobState::Enum state = sampleState(sample);
351
352 sample->meta()->setString("nc_ELG_state_details", "");
353
354 Status::Enum status = Status::PENDING;
355 switch (state) {
356 case JobState::INIT:
357 status = submit(sample, isFirstSample);
358 break;
359 case JobState::RUN:
360 status = checkPandaTask(sample);
361 break;
362 case JobState::DOWNLOAD:
363 status = download(sample);
364 break;
365 case JobState::MERGE:
366 status = merge(sample);
367 break;
368 case JobState::FINISHED:
369 case JobState::FAILED:
370 break;
371 }
372
373 state = nextState(state, status);
374 sample->meta()->setString("nc_ELG_state", JobState::name[state]);
375}
static Status::Enum submit(SH::Sample *const sample, const bool isFirstSample)
static JobState::Enum nextState(JobState::Enum state, Status::Enum status)
static Status::Enum checkPandaTask(SH::Sample *const sample)
static Status::Enum download(SH::Sample *const sample)
Definition merge.py:1

◆ sampleState()

JobState::Enum sampleState ( SH::Sample * sample)
static

Definition at line 101 of file PrunDriver.cxx.

102{
103 RCU_REQUIRE(sample);
104 static const std::string defaultState = JobState::name[JobState::INIT];
105 std::string label = sample->meta()->castString ("nc_ELG_state", defaultState, SH::MetaObject::CAST_NOCAST_DEFAULT);
106 return JobState::parse(label);
107}
std::string label(const std::string &format, int i)
Definition label.h:19

◆ saveJobDef()

void saveJobDef ( const std::string & fileName,
const EL::Job & job,
const SH::SampleHandler sh )
static

Definition at line 458 of file PrunDriver.cxx.

461{
462 TFile file(fileName.c_str(), "RECREATE");
463 TList outputs;
464 for (EL::Job::outputIter o = job.outputBegin(); o !=job.outputEnd(); ++o)
465 outputs.Add(o->Clone());
466 file.WriteTObject(&job.jobConfig(), "jobConfig", "SingleKey");
467 file.WriteTObject(&outputs, "outputs", "SingleKey");
468 bool haveDefault = false;
469 for (SH::SampleHandler::iterator s = sh.begin(); s != sh.end(); ++s) {
470 const SH::MetaObject& meta = *((*s)->meta());
471 file.WriteObject(&meta, meta.castString("sample_name").c_str());
472 if (!haveDefault)
473 {
474 file.WriteObject (&meta, "defaultMetaObject");
475 haveDefault = true;
476 }
477 }
478}
-diff
TFile * file

◆ submit()

Status::Enum submit ( SH::Sample *const sample,
const bool isFirstSample )
static

Definition at line 179 of file PrunDriver.cxx.

180{
181 RCU_REQUIRE(sample);
182 using namespace EL::msgEventLoop;
183
184 ANA_MSG_INFO( "Submitting " << sample->name() << "..." );
185
186 static bool loaded = false;
187 if (not loaded) {
188 // TString path = "$ROOTCOREBIN/python/EventLoopGrid/ELG_prun.py";
189 // gSystem->ExpandPathName(path);
190 // TPython::LoadMacro(path.Data());
191 std::string path = PathResolverFindCalibFile("EventLoopGrid/ELG_prun.py");
192 TPython::LoadMacro(path.c_str());
193 loaded = true;
194 }
195
196 TPython::Bind(dynamic_cast<TObject*>(sample), "ELG_SAMPLE");
197#if ROOT_VERSION_CODE >= ROOT_VERSION(6,33,01)
198 std::any result;
199 TPython::Exec("_anyresult = ROOT.std.make_any['int'](ELG_prun(ELG_SAMPLE))", &result);
200 int ret = std::any_cast<int>(result);
201#else
202 int ret = TPython::Eval("ELG_prun(ELG_SAMPLE)");
203#endif
204 TPython::Bind(0, "ELG_SAMPLE");
205
206 // Tarball is created for the first sample to be submitted
207 // then the tarball is simply reused for the other samples
208 // If the returned value is 1 it implies the tarball creation failed
209 // See EventLoopGrid/data/ELG_prun.py script
210 // Abort any further processing as the tarball was not succesfully created
211 if (isFirstSample && ret == 1){
212 ANA_MSG_ERROR("Failed to create tarball");
213 throw std::runtime_error("PrunDriver.cxx: aborting due to tarball creation issue");
214 }
215
216 if (ret < 100) {
217 sample->meta()->setString("nc_ELG_state_details",
218 "problem submitting");
219 return Status::FAIL;
220 }
221
222 sample->meta()->setDouble("nc_jediTaskID", ret);
223
224 // Let's also tell people about their task ID
225 ANA_MSG_INFO( "Task submitted; jediTaskID=" << ret );
226
227 return Status::DONE;
228}
#define ANA_MSG_ERROR(xmsg,...)
Macro printing error messages.
#define ANA_MSG_INFO(xmsg,...)
Macro printing info messages.