import logging import logging.handlers import time import json import pika import subprocess import sys import threading import os from distutils.util import strtobool import shlex class JixelGeneralPurposeRabbitConsumer(): def __init__( self, rabbitmq_username, rabbitmq_password, rabbitmq_host, rabbitmq_vhost, rabbitmq_port, rabbitmq_queue, rabbitmq_heartbeat, execution_path, execution_command, logger, discarding_mode=False ): self.rabbitmq_username = rabbitmq_username self.rabbitmq_password = rabbitmq_password self.rabbitmq_host = rabbitmq_host self.rabbitmq_vhost = rabbitmq_vhost self.rabbitmq_port = rabbitmq_port self.rabbitmq_queue = rabbitmq_queue self.rabbitmq_heartbeat = int(rabbitmq_heartbeat) self.execution_path = execution_path self.execution_command = execution_command self.logger = logger self.discarding_mode = strtobool(discarding_mode) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - Initialization...") self.credentials = pika.PlainCredentials( self.rabbitmq_username , self.rabbitmq_password ) self.connection = pika.BlockingConnection( pika.ConnectionParameters( host=self.rabbitmq_host, virtual_host=self.rabbitmq_vhost, port=int(self.rabbitmq_port), credentials=self.credentials,heartbeat=self.rabbitmq_heartbeat ) ) self.channel = self.connection.channel() self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - Initialization completed successfully") self.channel.queue_declare(queue=self.rabbitmq_queue, durable=True) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - " + self.rabbitmq_queue + " queue declared") self.channel.queue_declare(queue=self.rabbitmq_queue+'_errors', durable=True) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - " + self.rabbitmq_queue + "_errors queue declared") self.channel.basic_qos(prefetch_count=1) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - 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(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + ' - Processing thread is running:') string_command = "{}cake {} '{}'".format(self.execution_path, self.execution_command, json.dumps(mr)) command = shlex.split(string_command) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + ' - Executing shell command:') self.logger.info(command) try: completed_process = subprocess.run(command, stdout=subprocess.PIPE, stderr=subprocess.PIPE) if completed_process.stderr: self.logger.error(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " -> " + self.execution_command + ':\n{}'.format(completed_process.stderr.decode('utf-8'))) if completed_process.stdout: self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " -> " + self.execution_command + ':\n{}'.format(completed_process.stdout.decode('utf-8'))) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " -> " + self.execution_command + ' - RETURN CODE: {}'.format(completed_process.returncode)) completed_process.check_returncode() self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " -> " + self.execution_command + ': Shell command successfully executed!') self.result = True except FileNotFoundError as error: self.logger.exception(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " -> " + self.execution_command + ' - A FileNotFoundError error has occured during shell command execution:', stack_info=False, exc_info=False) self.logger.exception(error, stack_info=False, exc_info=False) return False except subprocess.CalledProcessError as error: self.logger.exception(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " -> " + self.execution_command + ' - A CalledProcessError has occured during processing:', stack_info=False, exc_info=False) self.logger.exception(error, stack_info=False, exc_info=False) self.result = False def __data_processing(self, channel, mr): self.logger.info(mr) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + ' - 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(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - 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): channel.basic_ack(delivery_tag=method.delivery_tag) if ack is True: self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - Ack performed") else: self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - Ack performed BUT ERRORS OCCURRED!") def __data_handler(self, channel, method, properties, body): self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + ' - 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.exception(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - Malformed message", stack_info=False, exc_info=False) self.logger.exception(error, stack_info=False, exc_info=False) self.logger.exception(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - This message will be removed from the queue", stack_info=False, exc_info=False) 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(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - An error has occured during processing") if self.discarding_mode is True: self.logger.error(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - Discarding mode enabled.") self.__send_ack_nack(True, channel, method) else: self.__send_ack_nack(False, channel, method) time.sleep(5) self.logger.info(self.rabbitmq_vhost + "/" + self.rabbitmq_queue + " - 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__': formatter = logging.Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s") handler = logging.handlers.TimedRotatingFileHandler("{}{}".format(os.environ['logging_file_path'], os.environ['logging_name']), when=os.environ['logging_file_rotate_when'], interval=int(os.environ['logging_file_rotate_interval']), backupCount=int(os.environ['logging_file_backup_count'])) handler.setFormatter(formatter) logger = logging.getLogger(os.environ['logging_name']) logger.addHandler(handler) logger.setLevel(logging.DEBUG) app = JixelGeneralPurposeRabbitConsumer( os.environ['RABBITMQ_USERNAME'], os.environ['RABBITMQ_PASSWORD'], os.environ['rabbitmq_host'], os.environ['rabbitmq_vhost'], os.environ['rabbitmq_port'], os.environ['rabbitmq_queue'], os.environ['rabbitmq_heartbeat'], os.environ['execution_path'], os.environ['execution_command'], logger, os.environ['discarding_mode'] ) app.consume()