MCPcopy Create free account
hub / github.com/diffgram/diffgram / process_new_event

Method process_new_event

eventhandlers/JobsConsumer.py:63–130  ·  view source on GitHub ↗

Receives a trigger event a fi queueclient = QueueClient()nds any actions matching the event trigger for action execution. :return:

(channel, method, properties, msg)

Source from the content-addressed store, hash-verified

61
62 @staticmethod
63 def process_new_event(channel, method, properties, msg):
64 """
65 Receives a trigger event a fi queueclient = QueueClient()nds any actions matching the event trigger for action
66 execution.
67 :return:
68 """
69 logger.debug(f'New Job Message: {msg}')
70 msg_data = json.loads(msg)
71 log = regular_log.default()
72 if msg_data.get('task_template_id') is None:
73 log['error']['task_template_id'] = f'Message most contain task template ID. Message is: {msg_data}'
74 if msg_data.get('file_id_list') is None:
75 log['error']['file_id_list'] = f'Message most contain a file_id_list. Message is: {msg_data}'
76 if msg_data.get('member_id') is None:
77 log['error']['member_id'] = f'Message most contain a member_id. Message is: {msg_data}'
78
79 if regular_log.log_has_error(log):
80 logger.error(f'Error processing jobs message: {log}')
81 return
82
83 task_template_id = msg_data.get('task_template_id')
84 file_id_list = msg_data.get('file_id_list')
85 member_id = msg_data.get('member_id')
86 logger.debug(f'Creating tasks for Job: {msg}')
87
88 with session_scope_threaded() as session:
89 task_template = Job.get_by_id(session = session, job_id = task_template_id)
90 files = File.get_by_id_list(session = session, file_id_list = file_id_list)
91 member = Member.get_by_id(session = session, member_id = member_id)
92 job_sync_manager = job_dir_sync_utils.JobDirectorySyncManager(
93 session = session,
94 job = task_template,
95 log = log
96 )
97 for file in files:
98 directories_ids = File.get_directories_ids(session = session, file_id = file.id)
99 directory = WorkingDir.get_by_id(session = session, directory_id = directories_ids[0])
100 sync_event_manager = SyncEventManager.create_sync_event_and_manager(
101 session = session,
102 dataset_source_id = directory,
103 dataset_destination = None,
104 description = 'Sync file from dataset {} to job {} and create task'.format(
105 directory.nickname,
106 task_template.name
107 ),
108 file = file,
109 job = task_template,
110 input_id = file.input_id,
111 project = task_template.project,
112 event_effect_type = 'create_task',
113 event_trigger_type = 'file_added',
114 status = 'init',
115 member_created = member
116 )
117
118 task, log = job_sync_manager.add_file_into_job(
119 file = file,
120 incoming_directory = directory,

Callers

nothing calls this directly

Calls 8

add_file_into_jobMethod · 0.95
session_scope_threadedFunction · 0.90
get_by_id_listMethod · 0.80
get_directories_idsMethod · 0.80
getMethod · 0.45
get_by_idMethod · 0.45

Tested by

no test coverage detected