Non puoi selezionare più di 25 argomenti
Gli argomenti devono iniziare con una lettera o un numero, possono includere trattini ('-') e possono essere lunghi fino a 35 caratteri.
169 righe
7.4 KiB
169 righe
7.4 KiB
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()
|
|
|