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

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()