Skip to content

Commit e1a7e2e

Browse files
committed
adding simple engine only for CRAB workflow part
1 parent 93f6138 commit e1a7e2e

1 file changed

Lines changed: 111 additions & 0 deletions

File tree

Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
from weakref import WeakKeyDictionary
2+
import automation_control as ctrl
3+
import argparse
4+
import enum
5+
import logging
6+
from typing import Any, Type, Union
7+
from os import listdir, walk, environ
8+
from os.path import isfile, join
9+
10+
logger = logging.getLogger("EfficiencyAnalysisLogger")
11+
logger.setLevel(logging.DEBUG)
12+
13+
ch = logging.FileHandler("EfficiencyAnalysisEngine.log")
14+
ch.setLevel(logging.DEBUG)
15+
formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
16+
ch.setFormatter(formatter)
17+
logger.addHandler(ch)
18+
19+
ch = logging.StreamHandler()
20+
ch.setLevel(logging.INFO)
21+
formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
22+
ch.setFormatter(formatter)
23+
logger.addHandler(ch)
24+
25+
campaign=environ.get("CAMPAIGN")
26+
workflow=environ.get("WORKFLOW")
27+
dataset=environ.get("DATASET")
28+
proxy=environ.get("PROXY")
29+
30+
template_for_first_module = "TempSteps/CrabConfigTemplateForFirstModule.py"
31+
32+
@ctrl.define_status_enum
33+
class TaskStatusEnum(enum.Enum):
34+
"""
35+
Class to encode enum tasks statuses for the purpouse of this automation workflow
36+
"""
37+
initialized = enum.auto(),
38+
duringFirstWorker = enum.auto(),
39+
waitingForFirstWorkerTransfer= enum.auto()
40+
done = enum.auto()
41+
42+
43+
@ctrl.decorate_with_enum(TaskStatusEnum)
44+
class TaskStatus:
45+
loop_id = 0.0
46+
condor_job_id = 0
47+
48+
def get_tasks_numbers_list(tasks_list_path):
49+
with open(tasks_list_path) as tasks_list_path:
50+
tasks_list_data = tasks_list_path.read()
51+
tasks_list_data = tasks_list_data.replace(" ", "")
52+
tasks_list = tasks_list_data.split(",")
53+
return tasks_list
54+
55+
56+
def prepare_parser()->argparse.ArgumentParser:
57+
parser = argparse.ArgumentParser(description=
58+
"""This is a script to run PPS Efficiency Analysis automation workflow""", formatter_class=argparse.RawTextHelpFormatter)
59+
60+
parser.add_argument('-t', '--tasks_list', dest='tasks_list_path', help='path to file containing list of data periods', required=True)
61+
return parser
62+
63+
64+
def get_runs_range(data_period):
65+
"""MOCKED"""
66+
return '317080'
67+
68+
69+
def process_new_tasks(tasks_list_path, task_controller):
70+
tasks_list = get_tasks_numbers_list(tasks_list_path)
71+
tasks_list = set(tasks_list)
72+
tasks_in_database = task_controller.getAllTasks().get_points()
73+
tasks_in_database = set(map(lambda x: x['dataPeriod'], tasks_in_database))
74+
tasks_not_submited_yet = tasks_list-tasks_in_database
75+
if tasks_not_submited_yet:
76+
task_controller.submitTasks(tasks_not_submited_yet)
77+
78+
79+
def submit_task_to_crab(campaign, workflow, data_period, dataset, template, proxy):
80+
result = ctrl.submit_task_to_crab(campaign, workflow, data_period, get_runs_range(data_period), template, dataset, proxy)
81+
82+
return result
83+
84+
85+
def set_status_after_first_worker_submission(task_status, operation_result):
86+
task_status.duringFirstWorker=1
87+
task_status.initialized=0
88+
task_status.loop_id+=1
89+
return task_status
90+
91+
92+
storage_path = "/eos/user/m/mobrzut"
93+
94+
TRANSITIONS_DICT = {
95+
'initialized': (submit_task_to_crab, 0, set_status_after_first_worker_submission, [dataset, template_for_first_module, proxy] ),
96+
'duringFirstWorker': (ctrl.check_if_crab_task_is_finished, True, TaskStatus.waitingForFirstWorkerTransfer, [proxy]),
97+
'waitingForFirstWorkerTransfer': (ctrl.is_crab_output_already_transfered, True, TaskStatus.duringFirstHarvester, [proxy])
98+
}
99+
100+
101+
102+
if __name__ == '__main__':
103+
parser = prepare_parser()
104+
opts = parser.parse_args()
105+
task_controller = ctrl.TaskCtrl.TaskControl(campaign=campaign, workflow=workflow, TaskStatusClass=TaskStatus)
106+
process_new_tasks(opts.tasks_list_path, task_controller)
107+
finite_state_machine = ctrl.FiniteStateMachine(TRANSITIONS_DICT)
108+
finite_state_machine.process_tasks(task_controller, TaskStatusClass=TaskStatus)
109+
110+
111+

0 commit comments

Comments
 (0)