52 void fillJob (BatchJob& myjob,
const Job& job,
53 const std::string& submitDir)
56 myjob.location = submitDir;
66 void fillSample (BatchSample& mysample,
69 mysample.name = sample.name();
70 mysample.meta = *sample.meta();
71 mysample.meta.fetchDefaults (
meta);
73 mysample.files = sample.makeFileList ();
84 void addSample (BatchJob& job, BatchSample sample,
85 const std::vector<BatchSegment>& segments)
89 sample.begin_segments = job.segments.size();
93 for (std::size_t iter = 0, end = segments.size(); iter != end; ++ iter)
95 BatchSegment segment = segments[iter];
97 segment.sampleName = sample.name;
99 std::ostringstream myname;
100 myname << sample.name <<
"-" << iter;
101 segment.fullName = myname.str();
104 std::ostringstream myname;
106 segment.segmentName = myname.str();
109 segment.sample = job.samples.size();
110 segment.job_id = job.segments.size();
113 segment.end_file = segments[iter+1].begin_file;
114 segment.end_event = segments[iter+1].begin_event;
117 segment.end_file = sample.files.size();
118 segment.end_event = 0;
120 RCU_ASSERT (segment.begin_file < segment.end_file || (segment.begin_file == segment.end_file && segment.begin_event <= segment.end_event));
121 job.segments.push_back (segment);
124 sample.end_segments = job.segments.size();
125 job.samples.push_back (sample);
139 void splitSampleByFile (
const BatchSample& sample,
140 std::vector<BatchSegment>& segments,
147 for (std::size_t
index = 0, end = numJobs;
150 BatchSegment segment;
151 segment.begin_file = rint (
index * (sample.files.size() / numJobs));
152 segments.push_back (segment);
167 void splitSampleByEvent (
const BatchSample& sample,
168 std::vector<BatchSegment>& segments,
169 const std::vector<Long64_t>& eventsFile,
173 RCU_REQUIRE (sample.files.size() == eventsFile.size());
176 Long64_t eventsSum = 0;
177 for (std::vector<Long64_t>::const_iterator nevents = eventsFile.begin(),
178 end = eventsFile.end(); nevents != end; ++ nevents)
181 eventsSum += *nevents;
189 segments.push_back (BatchSegment ());
192 if (numJobs > eventsSum)
197 numJobs = ceil (
double (eventsSum) / eventsMax);
198 eventsMax = Long64_t (ceil (
double (eventsSum) / numJobs));
200 BatchSegment segment;
201 while (std::size_t (segment.begin_file) < eventsFile.size())
203 while (segment.begin_event < eventsFile[segment.begin_file])
205 segments.push_back (segment);
206 segment.begin_event += eventsMax;
208 segment.begin_event -= eventsFile[segment.begin_file];
209 ++ segment.begin_file;
227 void splitSample (
const BatchSample& sample,
228 std::vector<BatchSegment>& segments)
237 const double filesPerWorker
239 if (filesPerWorker < 1)
241 std::ostringstream
msg;
242 msg <<
"invalid number of files per worker: " << filesPerWorker;
243 throw std::runtime_error (
msg.str());
245 double numJobs = ceil (sample.files.size() / filesPerWorker);
255 splitSampleByEvent (sample, segments, meta_nentries->
value, numJobs);
258 splitSampleByFile (sample, segments, numJobs);
262 if (segments.empty())
267 segments.push_back (
empty);
278 void fillFullJob (BatchJob& myjob,
const Job& job,
279 const std::string& location,
282 fillJob (myjob, job, location);
283 *myjob.job.options() =
meta;
285 for (
const SH::Sample *sample : job.sampleHandler())
287 BatchSample mysample;
288 fillSample (mysample, *sample,
meta);
290 std::vector<BatchSegment> subsegments;
291 splitSample (mysample, subsegments);
293 addSample (myjob, mysample, subsegments);
301 testInvariant ()
const
315 doManagerStep (Detail::ManagerData& data)
const
317 using namespace msgEventLoop;
323 data.sharedFileSystem
330 if (data.sharedFileSystem)
331 data.batchWriteLocation = data.submitDir;
333 data.batchWriteLocation =
".";
335 if (data.sharedFileSystem)
336 data.batchSubmitLocation = data.submitDir+
"/submit";
338 data.batchSubmitLocation =
".";
344 const std::string writeLocation=data.batchWriteLocation;
346 end = data.job->outputEnd(); out != end; ++ out)
348 if (out->output() ==
nullptr)
351 (writeLocation +
"/fetch/data-" + out->label() +
"/"));
359 ANA_MSG_DEBUG (
"submitting batch job in location " << data.submitDir);
360 const std::string submitDir = data.submitDir +
"/submit";
361 if (gSystem->MakeDirectory (submitDir.c_str()) != 0)
364 return ::StatusCode::FAILURE;
366 const std::string runDir = data.submitDir +
"/run";
367 if (gSystem->MakeDirectory (runDir.c_str()) != 0)
370 return ::StatusCode::FAILURE;
372 const std::string fetchDir = data.submitDir +
"/fetch";
373 if (gSystem->MakeDirectory (fetchDir.c_str()) != 0)
376 return ::StatusCode::FAILURE;
378 const std::string statusDir = data.submitDir +
"/status";
379 if (gSystem->MakeDirectory (statusDir.c_str()) != 0)
382 return ::StatusCode::FAILURE;
389 data.batchJob = std::make_unique<BatchJob> ();
390 fillFullJob (*data.batchJob, *data.job, data.submitDir, data.options);
391 data.batchJob->location=data.batchWriteLocation;
393 std::string path = data.submitDir +
"/submit/config.root";
394 std::unique_ptr<TFile>
file (TFile::Open (path.c_str(),
"RECREATE"));
395 if (
file ==
nullptr ||
file->IsZombie())
398 return ::StatusCode::FAILURE;
403 if (
file->WriteObject (data.batchJob.get(),
"job") <= 0)
405 ANA_MSG_ERROR (
"failed to write job configuration to " << path);
406 return ::StatusCode::FAILURE;
410 std::ofstream
file ((data.submitDir +
"/submit/segments").c_str());
411 for (std::size_t iter = 0, end = data.batchJob->segments.size();
412 iter != end; ++ iter)
414 file << iter <<
" " << data.batchJob->segments[iter].fullName <<
"\n";
422 data.batchName =
"run";
424 data.batchJobId =
"EL_JOBID=$1\n";
430 makeScript (data, data.batchJob->segments.size());
436 for (std::size_t
index = 0;
index != data.batchJob->segments.size(); ++
index)
437 data.batchJobIndices.push_back (
index);
444 ANA_MSG_INFO (
"retrieving batch job in location " << data.submitDir);
446 std::unique_ptr<TFile>
file
447 {TFile::Open ((data.submitDir +
"/submit/config.root").c_str(),
"READ")};
448 if (
file ==
nullptr ||
file->IsZombie())
451 return ::StatusCode::FAILURE;
453 data.batchJob.reset (
dynamic_cast<BatchJob*
>(
file->Get (
"job")));
454 if (data.batchJob ==
nullptr)
457 return ::StatusCode::FAILURE;
459 data.job = &data.batchJob->job;
466 for (std::size_t job = 0; job != data.batchJob->segments.size(); ++ job)
468 std::ostringstream completedFile;
469 completedFile << data.submitDir <<
"/status/completed-" << job;
470 const bool hasCompleted =
471 (gSystem->AccessPathName (completedFile.str().c_str()) == 0);
473 std::ostringstream failFile;
474 failFile << data.submitDir <<
"/status/fail-" << job;
476 (gSystem->AccessPathName (failFile.str().c_str()) == 0);
482 ANA_MSG_ERROR (
"sub-job " << job <<
" reported both success and failure");
483 return ::StatusCode::FAILURE;
485 data.batchJobSuccess.insert (job);
489 data.batchJobFailure.insert (job);
491 data.batchJobUnknown.insert (job);
494 ANA_MSG_INFO (
"current job status: " << data.batchJobSuccess.size() <<
" success, " << data.batchJobFailure.size() <<
" failure, " << data.batchJobUnknown.size() <<
" running/unknown");
500 bool all_missing =
false;
501 if (data.resubmitOption ==
"ALL_MISSING")
504 }
else if (!data.resubmitOption.empty())
506 ANA_MSG_ERROR (
"unknown resubmit option " + data.resubmitOption);
507 return ::StatusCode::FAILURE;
510 for (std::size_t segment = 0; segment != data.batchJob->segments.size(); ++ segment)
514 if (data.batchJobSuccess.find (segment) == data.batchJobSuccess.end())
515 data.batchJobIndices.push_back (segment);
518 if (data.batchJobFailure.find (segment) != data.batchJobFailure.end())
519 data.batchJobIndices.push_back (segment);
523 if (data.batchJobIndices.empty())
527 return ::StatusCode::SUCCESS;
530 for (std::size_t segment : data.batchJobIndices)
532 std::ostringstream command;
534 command <<
" " <<
RCU::Shell::quote (data.submitDir) <<
"/status/completed-" << segment;
535 command <<
" " <<
RCU::Shell::quote (data.submitDir) <<
"/status/fail-" << segment;
536 command <<
" " <<
RCU::Shell::quote (data.submitDir) <<
"/status/done-" << segment;
539 data.options = *data.batchJob->job.options();
545 bool merged = mergeHists (data);
548 diskOutputSave (data);
550 data.retrieved =
true;
551 data.completed = merged;
558 return ::StatusCode::SUCCESS;
563 std::string BatchDriver ::
564 defaultReleaseSetup (
const Detail::ManagerData& data)
const
569 const std::string tarballName(
"AnalysisPackage.tar.gz");
571 std::ostringstream
file;
574 const char *WORKDIR_DIR = getenv (
"WorkDir_DIR");
577 std::string CMAKE_DIR_str;
578 if (WORKDIR_DIR ==
nullptr){
579 msgEventLoop::ANA_MSG_INFO (
"Could not find environment variable $WorkDir_DIR");
580 const char *CMAKE_PREFIX_PATH = getenv (
"CMAKE_PREFIX_PATH");
581 if (CMAKE_PREFIX_PATH ==
nullptr){
582 throw std::runtime_error (
"neither $WorkDir_DIR nor $CMAKE_PREFIX_PATH is set, cannot locate the release");
585 CMAKE_DIR_str = CMAKE_PREFIX_PATH;
586 if (CMAKE_DIR_str.find(
":") != std::string::npos){
588 CMAKE_DIR_str.erase( CMAKE_DIR_str.find(
":") , std::string::npos );
591 WORKDIR_DIR = CMAKE_DIR_str.data();
594 if(!data.sharedFileSystem)
596 file <<
"mkdir -p build && tar -C build/ -xf " << tarballName <<
" || abortJob\n";
603 if(getenv(
"AtlasSetupSite"))
file <<
"export AtlasSetupSite=" << getenv(
"AtlasSetupSite") <<
"\n";
605 if(getenv(
"AtlasSetup"))
file <<
"export AtlasSetup=" << getenv(
"AtlasSetup") <<
"\n";
609 file <<
"export ATLAS_LOCAL_ROOT_BASE=/cvmfs/atlas.cern.ch/repo/ATLASLocalRootBase\n";
610 file <<
"source $ATLAS_LOCAL_ROOT_BASE/user/atlasLocalSetup.sh --quiet\n";
613 std::ostringstream defaultSetupCommand;
616 if(getenv(
"AtlasProject")) defaultSetupCommand <<
"export AtlasProject=" << getenv(
"AtlasProject") <<
"\n";
618 if(getenv(
"AtlasVersion")) defaultSetupCommand <<
"export AtlasVersion=" << getenv(
"AtlasVersion") <<
"\n";
620 if(getenv(
"AtlasBuildStamp")) defaultSetupCommand <<
"export AtlasBuildStamp=" << getenv(
"AtlasBuildStamp") <<
"\n";
622 if(getenv(
"AtlasBuildBranch")) defaultSetupCommand <<
"export AtlasBuildBranch=" << getenv(
"AtlasBuildBranch") <<
"\n";
624 if(getenv(
"AtlasReleaseType")) defaultSetupCommand <<
"export AtlasReleaseType=" << getenv(
"AtlasReleaseType") <<
"\n";
625 defaultSetupCommand <<
"if [ \"${AtlasReleaseType}\" == \"stable\" ]; then\n";
626 defaultSetupCommand <<
" source ${AtlasSetup}/scripts/asetup.sh ${AtlasProject},${AtlasVersion} || abortJob\n";
627 defaultSetupCommand <<
"else\n";
628 defaultSetupCommand <<
" source ${AtlasSetup}/scripts/asetup.sh ${AtlasProject},${AtlasBuildBranch},${AtlasBuildStamp} || abortJob\n";
629 defaultSetupCommand <<
"fi\n";
630 defaultSetupCommand <<
"echo \"Using default setup command\"";
634 if(data.sharedFileSystem)
file <<
"source " << WORKDIR_DIR <<
"/setup.sh || abortJob\n";
635 else file <<
"source build/setup.sh || abortJob\n";
638 if(!data.sharedFileSystem)
640 std::ostringstream cmd;
642 cmd <<
"cpack -D CPACK_INSTALL_PREFIX=. -G TGZ --config $TestArea/CPackConfig.cmake";
645 if (gSystem->Exec (cmd.str().c_str()) != 0){
646 throw std::runtime_error (
"failed to execute: " + cmd.str());
649 const char *WorkDir_VERSION = std::getenv (
"WorkDir_VERSION");
650 const char *WorkDir_PLATFORM = std::getenv (
"WorkDir_PLATFORM");
651 if (WorkDir_VERSION ==
nullptr || WorkDir_PLATFORM ==
nullptr){
652 throw std::runtime_error (
"$WorkDir_VERSION and $WorkDir_PLATFORM must both be set to build the release tarball");
654 std::ostringstream mv_command;
655 mv_command <<
"mv WorkDir_" << WorkDir_VERSION <<
"_" << WorkDir_PLATFORM <<
".tar.gz " << tarballName;
656 if (gSystem->Exec (mv_command.str().c_str()) != 0){
657 throw std::runtime_error (
"failed to execute: " + mv_command.str());
667 makeScript (Detail::ManagerData& data,
668 std::size_t njobs)
const
672 std::string name = data.batchName;
673 bool multiFile = (name.find (
"{JOBID}") != std::string::npos);
679 const std::string releaseSetup =
680 data.batchSkipReleaseSetup ? std::string() : defaultReleaseSetup (data);
682 for (std::size_t
index = 0, end = multiFile ? njobs : 1;
index != end; ++
index)
684 std::ostringstream
str;
686 const std::string fileName = data.submitDir +
"/submit/" +
RCU::substitute (name,
"{JOBID}",
str.str());
689 std::ofstream
file (fileName.c_str());
690 file <<
"#!/bin/bash\n";
691 file <<
"echo starting batch job initialization\n";
693 file <<
"echo batch job user initialization finished\n";
694 if (multiFile)
file <<
"EL_JOBID=" <<
index <<
"\n\n";
695 else file << data.batchJobId <<
"\n";
697 file <<
"function abortJob {\n";
698 file <<
" echo \"abort EL_JOBID=${EL_JOBID}\"\n";
699 file <<
" touch \"" << data.batchWriteLocation <<
"/status/fail-$EL_JOBID\"\n";
700 file <<
" touch \"" << data.batchWriteLocation <<
"/status/done-$EL_JOBID\"\n";
705 file <<
"EL_JOBSEG=`grep \"^$EL_JOBID \" \"" << data.batchSubmitLocation <<
"/segments\" | awk ' { print $2 }'`\n";
706 file <<
"test \"$EL_JOBSEG\" != \"\" || abortJob\n";
707 file <<
"hostname\n";
710 file << shellInit <<
"\n";
712 if(!data.sharedFileSystem)
714 file <<
"mkdir \"fetch\" || abortJob\n";
715 file <<
"mkdir \"status\" || abortJob\n";
719 if(data.sharedFileSystem)
721 file <<
"test \"$TMPDIR\" == \"\" && TMPDIR=/tmp\n";
722 file <<
"RUNDIR=${TMPDIR}/EventLoop-Worker-$EL_JOBSEG-`date +%s`-$$\n";
723 file <<
"mkdir \"$RUNDIR\" || abortJob\n";
724 file <<
"cd \"$RUNDIR\" || abortJob\n";
727 if (!data.batchSkipReleaseSetup)
728 file << releaseSetup;
730 file <<
"eventloop_batch_worker $EL_JOBID '" << data.batchSubmitLocation <<
"/config.root' || abortJob\n";
732 file <<
"test -f \"" << data.batchWriteLocation <<
"/status/completed-$EL_JOBID\" || "
733 <<
"touch \"" << data.batchWriteLocation <<
"/status/fail-$EL_JOBID\"\n";
734 file <<
"touch \"" << data.batchWriteLocation <<
"/status/done-$EL_JOBID\"\n";
735 if(data.sharedFileSystem)
file <<
"cd .. && rm -rf \"$RUNDIR\"\n";
739 std::ostringstream cmd;
741 if (gSystem->Exec (cmd.str().c_str()) != 0)
742 throw std::runtime_error (
"failed to execute: " + cmd.str());
750 mergeHists (Detail::ManagerData& data)
752 using namespace msgEventLoop;
759 std::unique_ptr<SH::DiskOutput> origHistOutputMemory;
761 for (
auto iter = data.job->outputBegin(), end = data.job->outputEnd();
762 iter != end; ++ iter)
765 origHistOutput = iter->output();
767 if (origHistOutput ==
nullptr)
769 origHistOutputMemory = std::make_unique<SH::DiskOutputLocal>
770 (data.submitDir +
"/fetch/hist-");
771 origHistOutput = origHistOutputMemory.get();
777 ANA_MSG_DEBUG (
"merging histograms in location " << data.submitDir);
779 for (std::size_t sample = 0, end = data.batchJob->samples.size();
780 sample != end; ++ sample)
782 const BatchSample& mysample (data.batchJob->samples[sample]);
784 std::ostringstream output;
785 output << data.submitDir <<
"/hist-" << data.batchJob->samples[sample].name <<
".root";
786 if (gSystem->AccessPathName (output.str().c_str()) != 0)
788 ANA_MSG_VERBOSE (
"merge files for sample " << data.batchJob->samples[sample].name);
790 bool complete =
true;
791 std::vector<std::string> input;
792 for (std::size_t segment = mysample.begin_segments,
793 end = mysample.end_segments; segment != end; ++ segment)
795 const BatchSegment& mysegment = data.batchJob->segments[segment];
797 const std::string hist_file = origHistOutput->
targetURL
798 (mysegment.sampleName, mysegment.segmentName,
".root");
800 ANA_MSG_VERBOSE (
"merge segment " << segment <<
" completed=" << (data.batchJobSuccess.find(segment)!=data.batchJobSuccess.end()) <<
" fail=" << (data.batchJobFailure.find(segment)!=data.batchJobFailure.end()) <<
" unknown=" << (data.batchJobUnknown.find(segment)!=data.batchJobUnknown.end()));
802 input.push_back (hist_file);
804 if (data.batchJobFailure.find(segment)!=data.batchJobFailure.end())
806 std::ostringstream message;
807 message <<
"subjob " << segment <<
"/" << mysegment.fullName
809 throw std::runtime_error (message.str());
811 else if (data.batchJobSuccess.find(segment)==data.batchJobSuccess.end())
812 complete =
false, result =
false;
820 end = data.batchJob->job.outputEnd(); out != end; ++ out)
823 output << data.submitDir <<
"/data-" << out->label();
825 if(gSystem->AccessPathName(output.str().c_str()))
826 gSystem->mkdir(output.str().c_str(),
true);
829 output <<
"/" << data.batchJob->samples[sample].name <<
".root";
831 std::vector<std::string> dataInput;
832 for (std::size_t segment = mysample.begin_segments,
833 dataEnd = mysample.end_segments; segment != dataEnd; ++ segment)
835 const BatchSegment& mysegment = data.batchJob->segments[segment];
837 const std::string infile =
838 data.submitDir +
"/fetch/data-" + out->label() +
"/" + mysegment.fullName +
".root";
840 dataInput.push_back (infile);