Merge commit '57598e815bd8554291636435905a1986c1b486f4' as 'utils_v2'

This commit is contained in:
2024-11-12 10:47:34 +05:30
89 changed files with 15120 additions and 0 deletions
+554
View File
@@ -0,0 +1,554 @@
"""
AUTHOR:
Khushal P Soonderji
DATE:
Saturday, 13th Jul, 2024
OBJECTIVE:
To be able to work with keys
REFERENCES:
N/A
DOWNLOADS:
N/A
"""
# *****************************************************************************************************************
# ***** ****
# *** IMPORT ***
# ***** ****
# *****************************************************************************************************************
# To make sibling directories accessible for imports:
import sys
sys.path.append(".")
sys.path.append("..")
# To use Kafka:
from aiokafka import AIOKafkaProducer
from aiokafka import AIOKafkaConsumer
# For working with JSON strings:
from utils_v2.string import json
from utils_v2.serialization.json_serializer import JSONSerializer
# For debugging:
from icecream import IceCreamDebugger
# For SSL security:
import ssl
# For asynchronous activities:
import asyncio
# *****************************************************************************************************************
# ***** ****
# *** MACROS / ONE-TIME INIT ***
# ***** ****
# *****************************************************************************************************************
# --- Nothing Yet
# *****************************************************************************************************************
# ***** ****
# *** VARIABLES ***
# ***** ****
# *****************************************************************************************************************
# --- Nothing Yet
# *****************************************************************************************************************
# ***** ****
# *** FUNCTIONS ***
# ***** ****
# *****************************************************************************************************************
def get_ssl_context(
ca_file,
cert_file,
key_file
):
"""
Generate the SSL context to use with the Kafka instances.
:param ca_file: The Certificate Authority file as a path to a local file.
:param cert_file: The Certificate file as a path to a local file.
:param key_file: The Key file as a path to a local file.
:return: The SSL context instance as a path to a local file.
"""
ssl_context = ssl.create_default_context()
ssl_context.load_verify_locations(ca_file)
ssl_context.load_cert_chain(certfile = cert_file, keyfile = key_file)
return ssl_context
# *****************************************************************************************************************
# ***** ****
# *** CLASSES ***
# ***** ****
# *****************************************************************************************************************
class ProducerKafka:
def __init__(
self,
topic,
serializer = None,
debug = True,
debug_prefix = "Kafka (P) | ",
**kwargs
):
"""
Create a Kafka Producer.
:param topic: The topic to produce on.
:param serializer: The serializer to use.
:param debug: Whether, or not, you want to print the debug strings.
:param debug_prefix: The prefix to use while debugging.
:param kwargs: Any configuration parameters for the Kafka instances.
"""
# Initialize the debugger:
self.__printer = IceCreamDebugger(prefix = debug_prefix, includeContext = True)
if not debug: self.__printer.disable()
# initialize the Kafka producer:
self.__topic = topic
self.__kwargs = kwargs
self.__producer = None
self.__connected = False
self.__serializer = serializer or JSONSerializer()
# For establishing connection:
self.__exclusive_semaphore = asyncio.Semaphore(1)
def enable_debug(self):
self.__printer.enable()
def disable_debug(self):
self.__printer.disable()
async def connect(self):
"""
Connects to the Kafka server if not connected.
:return: True or False based on the success of the operation.
"""
async with self.__exclusive_semaphore:
if not self.__connected:
try:
self.__producer = AIOKafkaProducer(**self.__kwargs)
await self.__producer.start()
self.__connected = True
except Exception as exception: self.__printer(exception)
return self.__connected
async def ensure_connection(self):
"""
Connects to the Kafka server if not connected.
:return: True or False based on the success of the operation.
"""
if not self.__connected: await self.connect()
return self.__connected
async def close(self):
"""
Terminates the connection.
:return: None.
"""
if self.__connected:
try:
await self.__producer.stop()
self.__printer("Producer closed!")
self.__connected = False
except Exception as exception: self.__printer(exception)
async def produce(self, value, key = None, topic = None, encoding = "utf-8"):
"""
Sends one message to the Kafka server on the topic that has been set for this instance.
:param value: The message to send.
:param key: The key to use when you want the messages to follow an order.
:param topic: A custom topic for this message, else the topic defined during the creation of this instance will
be used by default.
:param encoding: The encoding format.
:return: True or False based on the success of the operation.
"""
# Ensure connectivity to the server.
# If not connected, return with failure immediately.
if not await self.ensure_connection(): return False
try:
# Send the message:
await self.__producer.send_and_wait(
topic = topic or self.__topic,
value = self.__serializer.serialize(data = value, encoding = encoding),
key = key
)
# Return with success if no exception occurred:
return True
# Return with failure if something went wrong:
except Exception as exception:
self.__printer(exception, self.__topic, type(value), value)
return False
# ---------------------------------------------------------------------------------------------------------------------
class ConsumerKafka:
def __init__(
self,
topic,
serializer = None,
debug = True,
debug_prefix = "Kafka (C) | ",
**kwargs
):
"""
Create a Kafka Consumer.
:param topic: The topic to consumer on.
:param serializer: The serializer to use.
:param debug: Whether, or not, you want to print the debug strings.
:param debug_prefix: The prefix to use while debugging.
:param kwargs: Any configuration parameters for the Kafka instances.
"""
# Initialize the debugger:
self.__printer = IceCreamDebugger(prefix = debug_prefix, includeContext = True)
if not debug: self.__printer.disable()
# initialize the Kafka producer:
self.__topic = topic
self.__kwargs = kwargs
self.__consumer = None
self.__connected = False
self.__serializer = serializer or JSONSerializer()
# For establishing connection:
self.__exclusive_semaphore = asyncio.Semaphore(1)
def enable_debug(self):
self.__printer.enable()
def disable_debug(self):
self.__printer.disable()
async def connect(self):
"""
Connects to the Kafka server if not connected.
:return: True or False based on the success of the operation.
"""
async with self.__exclusive_semaphore:
if not self.__connected:
try:
self.__consumer = AIOKafkaConsumer(self.__topic, **self.__kwargs)
await self.__consumer.start()
self.__connected = True
except Exception as exception: self.__printer(exception)
return self.__connected
async def ensure_connection(self):
"""
Connects to the Kafka server if not connected.
:return: True or False based on the success of the operation.
"""
if not self.__connected: await self.connect()
return self.__connected
async def close(self):
"""
Terminates the connection.
:return: None.
"""
if self.__connected:
try:
await self.__consumer.stop()
self.__printer("Consumer closed!")
self.__connected = False
except Exception as exception: self.__printer(exception)
async def consume(self, count = 1, timeout = 0.05, encoding = "utf-8"):
"""
Get messages from the Kafka server.
:param count: The number of messages to get from the Kafka server.
:param timeout: The time in seconds to wait for retrieval.
:param encoding: The encoding to use.
:return: The messages that were received. If no messages are available, an empty list will be returned.
"""
# Ensure connectivity to the server.
# If not connected, return with failure immediately.
if not await self.ensure_connection(): return []
# Make a variable that will hold the final results:
messages = []
try:
# Read some messages:
results = await self.__consumer.getmany(
max_records = max(1, count),
timeout_ms = int(timeout * 1_000)
)
# Format the received messages:
if results:
for topic_partition, records in results.items():
for record in records:
record_dict = record.__dict__
record_dict["value"] = self.__serializer.deserialize(
data = record_dict["value"],
encoding = encoding
)
messages.append(record_dict)
# Debugging print if something went wrong:
except Exception as exception: self.__printer(exception)
# Done here:
return messages
# ---------------------------------------------------------------------------------------------------------------------
class BidirectionalKafka:
# The 'roles' that the instance can take.
# The master talks on the channel (topic) that the slave listens on and vice versa.
# Master-Slave is only for deciding who talks on which channel and who listens on which.
# In a two-party system, one must be the master, the other must be the slave.
# There are no extra privileges that the master enjoys. The naming convention was borrowed from common protocols
# used in electronics (like I2C).
ROLE_MASTER = 1
ROLE_SLAVE = 0
def __init__(
self,
role,
topic,
ack_topic: str = None,
group: str = None,
serializer = None,
debug = True,
debug_prefix = "Kafka (B) | ",
**kwargs
):
"""
Creates a walkie-talkie type setup to use Kafka in a bidirectional manner. Fo more information on all the
individual methods, please read through the doc-strings of the component classes 'ProducerKafka', and
'ConsumerKafka'.
:param role: Select from "ROLE_MASTER" and "ROLE_SLAVE". Between the two parties that are talking, one will be
the master and the other will be the slave. The channel that the master uses to speak will the one the slave
uses to listen, and vice versa.
:param topic: The topic to communicate on. Will be the same between the master and the slave.
:param ack_topic: Explicitly provide this for the second channel, or it will be created from the name of the
topic itself. Will be the same between the master and the slave.
:param group: The group to assign the instance to.
:param debug: Whether, or not, you want to print the debug strings.
:param debug_prefix: The prefix to use while debugging.
:param kwargs: Any configuration parameters for the Kafka instances.
"""
# Not down the basic variables:
self.__role = role
self.__topic = topic
self.__ack_topic = ack_topic or topic + "Ack"
self.__group = group
# In case the current instance is the master,
# it will talk on "topic", and listen on "ack_topic":
if self.__role == self.ROLE_MASTER:
self.__producer_kwargs = kwargs.copy()
self.__producer = ProducerKafka(
topic = self.__topic,
serializer = serializer,
debug = debug,
debug_prefix = debug_prefix.strip() + " (P) | ",
**self.__producer_kwargs
)
self.__consumer_kwargs = kwargs.copy()
self.__consumer_kwargs["group_id"] = self.__group
self.__consumer = ConsumerKafka(
topic = self.__ack_topic,
group = group,
serializer = serializer,
debug = debug,
debug_prefix = debug_prefix.strip() + " (C) | ",
**self.__consumer_kwargs
)
# On the other hand, if the current instance is a slave,
# It will listen on "topic", and talk on "ack_topic":
else:
self.__producer_kwargs = kwargs.copy()
self.__producer = ProducerKafka(
topic = self.__ack_topic,
serializer = serializer,
debug = debug,
debug_prefix = debug_prefix.strip() + " (P) | ",
**self.__producer_kwargs
)
self.__consumer_kwargs = kwargs.copy()
self.__consumer_kwargs["group_id"] = self.__group
self.__consumer = ConsumerKafka(
topic = self.__topic,
group = group,
serializer = serializer,
debug = debug,
debug_prefix = debug_prefix.strip() + " (C) | ",
**self.__consumer_kwargs
)
def enable_debug(self):
self.__producer.enable_debug()
self.__consumer.enable_debug()
def disable_debug(self):
self.__producer.disable_debug()
self.__consumer.disable_debug()
async def ensure_connection(self):
await self.__producer.ensure_connection()
await self.__consumer.ensure_connection()
async def close(self):
await self.__producer.close()
await self.__consumer.close()
async def produce(self, message, encoding = "utf-8"):
return await self.__producer.produce(message, encoding = encoding)
async def consume(self, count = 1, timeout = 0.05, encoding = "utf-8"):
return await self.__consumer.consume(count = count, timeout = timeout, encoding = encoding)
# *****************************************************************************************************************
# ***** ****
# *** MAIN PROGRAM ***
# ***** ****
# *****************************************************************************************************************
if __name__ == "__main__":
import asyncio
import time
from data_models.kafka_message import KafkaMessage
ssl_ctx = get_ssl_context(
ca_file = r"/home/developer/PycharmProjects/utils/cred/kafka/cert_authority.pem",
cert_file = r"/home/developer/PycharmProjects/utils/cred/kafka/fullchain.pem",
key_file = r"/home/developer/PycharmProjects/utils/cred/kafka/privkey.pem"
)
async def consumer_test():
consumer = ConsumerKafka(
topic = "kft_file_upload",
# group_id = "assessImg",
group_id = "updateMedia",
bootstrap_servers = "wtt.ditscentre.in:9092",
security_protocol = "SSL",
ssl_context = ssl_ctx
)
await consumer.connect()
await asyncio.sleep(1.5)
print("READY!")
while True:
messages = await consumer.consume(count = 1)
if len(messages) > 0: print("MESSAGE:", json.to_string(messages[0], default=str))
await asyncio.sleep(1.0)
async def producer_test():
producer = ProducerKafka(
topic = "kft_file_upload",
bootstrap_servers = "del.ditscentre.in:9092",
security_protocol = "SSL",
ssl_context = ssl_ctx
)
await producer.connect()
print("READY!")
while True:
my_msg = KafkaMessage(
data = {
"accepted": False,
"reason": "low resolution"
},
media = {
"name": "pikachu_poster.jpg",
"ext": "jpg",
"url": "https://nexcom.ditscentre.in/utils/files/small/download/66ded1c1c1c05139a618b5ff",
"attr": {
"user": "SarangKabir",
"project": "ACE-PGP",
"id": 173,
"campaignActivityId": "25",
"idCampaign": 49,
"phoneNo": "7977821877"
}
},
appId = "aceWockhardt",
proc = {
"name": "_assessImg",
"attr": {
"blurThreshold": 0.25,
"clarityThreshold": 0.65,
"nsfwThreshold": 0.25,
"minWidth": 512,
"minHeight": 512
}
},
ack = None
)
success = await producer.produce(my_msg.model_dump())
print("produced...")
time.sleep(1.0)
break
await producer.close()
asyncio.run(consumer_test())