Receives a trigger event a fi queueclient = QueueClient()nds any actions matching the event trigger for action execution. :return:
(channel, method, properties, msg)
| 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, |
nothing calls this directly
no test coverage detected