import logging import time import json import pika import subprocess import sys import threading import os from distutils.util import strtobool class JixelBackgroundTaskHandler(): def __init__( self, rabbitmq_username, rabbitmq_password, rabbitmq_host, rabbitmq_port, rabbitmq_queue, rabbitmq_heartbeat, execution_path, logger, discarding_mode=False ): self.rabbitmq_username = rabbitmq_username self.rabbitmq_password = rabbitmq_password self.rabbitmq_host = rabbitmq_host self.rabbitmq_port = rabbitmq_port self.rabbitmq_queue = rabbitmq_queue self.rabbitmq_heartbeat = int(rabbitmq_heartbeat) self.execution_path = execution_path self.logger = logger self.discarding_mode = strtobool(discarding_mode) self.logger.info("JixelBackgroundTaskHandler - Initialization...") self.credentials = pika.PlainCredentials( self.rabbitmq_username , self.rabbitmq_password ) self.connection = pika.BlockingConnection( pika.ConnectionParameters( host=self.rabbitmq_host, port=int(self.rabbitmq_port), credentials=self.credentials,heartbeat=self.rabbitmq_heartbeat ) ) self.channel = self.connection.channel() self.logger.info("JixelBackgroundTaskHandler - Initialization completed successfully") self.channel.queue_declare(queue=self.rabbitmq_queue, durable=True) self.logger.info("JixelBackgroundTaskHandler - " + self.rabbitmq_queue + " queue declared") self.channel.queue_declare(queue=self.rabbitmq_queue+'_errors', durable=True) self.logger.info("JixelBackgroundTaskHandler - " + self.rabbitmq_queue + "_errors queue declared") self.channel.basic_qos(prefetch_count=1) self.logger.info("JixelBackgroundTaskHandler - Queue handler declared") self.channel.basic_consume( on_message_callback=self.__data_handler, queue=self.rabbitmq_queue ) def consume(self): try: self.channel.start_consuming() except KeyboardInterrupt: self.channel.stop_consuming() self.channel.close() def __execute_command(self, mr): self.logger.info('JixelBackgroundTaskHandler - Processing thread is running:') command = [self.execution_path + 'cake'] command += mr['command'].split(' ') command += [mr['data']] self.logger.info('JixelBackgroundTaskHandler - Executing shell command:') self.logger.info(command) try: subprocess.check_call(command) self.result = True except subprocess.CalledProcessError as error: self.logger.error('JixelBackgroundTaskHandler - An error as occured during processing:') self.logger.error(error) self.result = False def __data_processing(self, channel, mr): self.logger.info(mr) self.logger.info('JixelBackgroundTaskHandler - Starting processing thread:') self.result = None thread = threading.Thread(target=self.__execute_command, args=(mr,)) thread.start() while thread.is_alive(): self.connection.process_data_events() self.logger.info("JixelBackgroundTaskHandler - Result received:") self.logger.info(self.result) if (self.result is not True): self.__send_error(json.dumps(mr)) return self.result def __send_ack_nack(self, ack, channel, method): if ack is True: self.logger.info("JixelBackgroundTaskHandler - Ack performed") channel.basic_ack(delivery_tag=method.delivery_tag) else: self.logger.info("JixelBackgroundTaskHandler - Nack performed BECAUSE ERRORS OCCURRED!") channel.basic_reject(delivery_tag=method.delivery_tag) def __data_handler(self, channel, method, properties, body): self.logger.info('JixelBackgroundTaskHandler - New message received:') error = False try: mr = json.loads(body.decode()) error = not self.__data_processing(channel, mr) except json.JSONDecodeError as error: self.logger.error("JixelBackgroundTaskHandler - Malformed message") self.logger.error(error) self.logger.error("JixelBackgroundTaskHandler - This message will be removed from the queue") self.__send_ack_nack(False, channel, method) time.sleep(60) return if self.result is True and error is False: self.__send_ack_nack(True, channel, method) else: self.logger.error("JixelBackgroundTaskHandler - An error has occured during processing") if self.discarding_mode is True: self.logger.error("JixelBackgroundTaskHandler - Discarding mode enabled.") self.__send_ack_nack(True, channel, method) else: self.__send_ack_nack(False, channel, method) time.sleep(5) self.logger.info("JixelBackgroundTaskHandler - Data handling complete. Waiting for new messages") def __send_error(self, message): self.channel.basic_publish(exchange='', routing_key=self.rabbitmq_queue+'_errors', body=message) if __name__ == '__main__': logger = logging.getLogger("JixelBackgroundTaskHandler") logger.setLevel(logging.INFO) formatter = logging.Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s") terminal_log = logging.StreamHandler(sys.stdout) terminal_log.setFormatter(formatter) logger.addHandler(terminal_log) app = JixelBackgroundTaskHandler( os.environ['RABBITMQ_USERNAME'], os.environ['RABBITMQ_PASSWORD'], os.environ['rabbitmq_host'], os.environ['rabbitmq_port'], os.environ['rabbitmq_queue'], os.environ['rabbitmq_heartbeat'], os.environ['execution_path'], logger, os.environ['discarding_mode'] ) app.consume()