From 30088dcea40d48320c4c56d668610db6d064f4c1 Mon Sep 17 00:00:00 2001 From: khushal Date: Mon, 9 Dec 2024 12:04:38 +0530 Subject: [PATCH] (20241209) Better session info capture in the logs. --- models/behaviour/core/__init__.py | 0 models/behaviour/core/auth_token.py | 221 +++++++++++++++++++ models/behaviour/mail/oauth_v3.py | 329 ++++++++++++++++++++++++++++ models/data/core/user_info.py | 212 ++++++++++++++++++ utils_v2/api/async_quart.py | 54 +++-- 5 files changed, 792 insertions(+), 24 deletions(-) create mode 100644 models/behaviour/core/__init__.py create mode 100644 models/behaviour/core/auth_token.py create mode 100644 models/behaviour/mail/oauth_v3.py create mode 100644 models/data/core/user_info.py diff --git a/models/behaviour/core/__init__.py b/models/behaviour/core/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/models/behaviour/core/auth_token.py b/models/behaviour/core/auth_token.py new file mode 100644 index 0000000..bb0bb99 --- /dev/null +++ b/models/behaviour/core/auth_token.py @@ -0,0 +1,221 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 5th Dec., 2024 + + OBJECTIVE: + + To create an interface between OpenAI and our internal system to perform LLM-based activities. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# My async utils: +from utils_v2.string import json +from utils_v2.date_time import date_time +from utils_v2.database.async_mysql_v2 import AsyncMySQL +from utils_v2.database.async_mongo_v2 import AsyncMongo + +# Base model: +from models.behaviour.base import BaseModel + +# Data Models: +from models.data.api.ai.llm import LLMInput, LLMOutput, LLMUsageTokens + +# To work with LLMs: +from langchain_openai import ChatOpenAI + +# To work with MongoDB: +from bson import ObjectId + +# To work with datatypes: +from typing import Literal + +# To make deep-copies: +import copy + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class LLMOpenAI(BaseModel): + + AI_USAGE_COLLECTION = "_aiUsage" + + def __init__( + self, + llm_creds: dict, + cache = None, + alert_url = None, + http_client = None, + debug = True, + debug_prefix = "Model | ", + debug_only_errors = True + ): + + """ + This is the model that works with OpenAi's LLM to perform tasks like text completion. + :param llm_creds: The JSON that holds the credentials to access your OpenAI account. Should have the keys + 'model', and 'openai_api_key'. + :param cache: The object to use for caching results from database calls. + :param alert_url: Which URL to call when something goes wrong. + :param http_client: The instance of an HTTP client to use when trying to send alerts and make other APIs. + :param debug: Whether, or not, you would like to print debugging messages: + :param debug_prefix: The prefix to print with the debugging messages. + :param debug_only_errors: Whether you would like to print only error messages or all messages. + :return: None. + """ + + # Initialize the parent: + super().__init__( + cache = cache, + alert_url = alert_url, + http_client = http_client, + debug = debug, + debug_prefix = debug_prefix, + debug_only_errors = debug_only_errors + ) + + # Create the interface to the LLM: + self.__llm = ChatOpenAI(**llm_creds) + + async def invoke( + self, + mongo_conn: AsyncMongo, + user_info: dict, + llm_input: LLMInput + ) -> LLMOutput: + + # Format the message as per the format of OpenAI: + prompt = [ + { + "role": {"system": "system", "ai": "assistant", "human": "user"}[message.role], + "content": message.content + } for message in llm_input.messages + ] + + # Invoke the AI, and format the response: + llm_response = await self.__llm.ainvoke(prompt) + llm_response = LLMOutput( + messages = llm_input.messages, + output = llm_response.content, + client = "openai", + model = llm_response.response_metadata["model_name"], + tokens = LLMUsageTokens( + input = llm_response.usage_metadata["input_tokens"], + output = llm_response.usage_metadata["output_tokens"], + total = llm_response.usage_metadata["total_tokens"], + ) + ) + + # Store this into MongoDB: + mongo_document = {"user": user_info} + for k, v in llm_response.model_dump().items(): mongo_document[k] = v + inserted_id = await mongo_conn.insert_one( + collection = self.AI_USAGE_COLLECTION, + document = mongo_document + ) + + # Done here: + return llm_response + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass + + # import asyncio + # + # llm_messages = [ + # { + # "role": "system", + # "content": "You are an office assistant." + # }, + # { + # "role": "ai", + # "content": "Hello, sir. How may I help you today?" + # }, + # { + # "role": "human", + # "content": "Please summarize this mail for me..." + # } + # ] + # + # my_llm = LLMOpenAI( + # llm_creds = { + # "model": "gpt-4o-mini", + # "openai_api_key": "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" + # } + # ) + # + # async def main(): + # + # llm_response = await my_llm.invoke(llm_input = LLMInput(messages = llm_messages)) + # print("LLM RESPONSE:", llm_response.model_dump_json(indent = 4)) + # + # asyncio.run(main()) diff --git a/models/behaviour/mail/oauth_v3.py b/models/behaviour/mail/oauth_v3.py new file mode 100644 index 0000000..04d66e9 --- /dev/null +++ b/models/behaviour/mail/oauth_v3.py @@ -0,0 +1,329 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Monday, 2nd Dec., 2024 + + OBJECTIVE: + + To define the interaction between the UI layer and the database connectivity in one place. Here we shall handle + all the activities for OAuth2.0 authorization requests for all the users of our service. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# My async utils: +from utils_v2.string import json +from utils_v2.date_time import date_time +from utils_v2.database.async_mysql_v2 import AsyncMySQL +from utils_v2.database.async_mongo_v2 import AsyncMongo + +# Base model: +from models.behaviour.base import BaseModel + +# To work with MongoDB: +from bson import ObjectId + +# To work with datatypes: +from typing import Literal + +# To make deep-copies: +import copy + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class MailOAuthModel(BaseModel): + + AUTH_COLLECTION = "_authTokens" + + async def get_token_id( + self, + db_conn: AsyncMySQL, + mongo_conn: AsyncMongo, + user_info: dict, + client_user_id: dict, + auth: dict, + service_client: Literal["gmail"], + auth_type: Literal["oauth"], + sync_freq: Literal[60, 300, 900] = 300, + session_token: str = None + ) -> ObjectId: + + """ + Stores params from the session info and gives an identifier to use in the authorization URL. Use this when the + user requests an authorization URL to link your service to another service (like GMail). + :param db_conn: The database connection (MariaDB) to use to perform the action. + :param mongo_conn: The database connection (MongoDB) to use to perform the action. + :param user_info: The dictionary that has the user's session information. + :param client_user_id: The way the third-party client recognizes your user. + :param auth: The authentication details of the account. + :param service_client: The name of the company or brand that is providing this service that is being integrated. + :param auth_type: To identify the type of authentication being done here. This could indicate simple password + authentication, more advance OAuth2.0 authentication, etc. + :param sync_freq: The time interval in which mails need to be sync'd. Specify this in seconds. + :param session_token: The session token of the user who requested this service. + :return: An ObjectId to later store the granted tokens. + """ + + # Note down the timestamp at which this event occurred: + request_ts = date_time.get_current_utc_date_time(as_string = False) + + # Get the identifier from the database: + mongo_json = await mongo_conn.find_one_and_update( + collection = MailOAuthModel.AUTH_COLLECTION, + filter = mongo_conn.dict_to_dot_notation({ + "serviceType": "email", + "user": { + "entityId": user_info["entityId"], + "billingAccountId": user_info["billingAccountId"] + }, + "clientUserId": client_user_id + }), + update = { + "$set": { + "lastRequestTs": request_ts, + "status": "active", + "syncFreq": max(sync_freq, 60) + }, + "$setOnInsert": { + "version": "1.1.1", + "serviceType": "email", + "client": service_client, + "authType": auth_type, + "user": user_info, + "clientUserId": client_user_id, + "auth": auth, + "token": None, + "firstRefreshTs": None, + "lastRefreshTs": None, + "firstRequestTs": request_ts, + } + }, + projection = { + "_id": True + }, + upsert = True, + return_updated = True + ) + + # Tell MariaDB that an authorization request was initiated: + db_json = {} + if mongo_json is not None: + db_json = await self.call_procedure( + db_conn = db_conn, + proc_name = "entity_integration_save", + proc_args = ( + user_info["entityId"], # ............................................ 'p_entity_id' + service_client, # ................................................... 'p_provider' + "Pending", # ........................................................ 'p_current_status' + "Auth Requested", # ................................................. 'p_last_action' + None, # ............................................................. 'p_display_name' + None, # ............................................................. 'p_display_picture' + str(mongo_json["_id"]), # ........................................... 'p_token_id' + json.to_string(python_data = {"email": None}, no_space = True), # ... 'p_notes' + user_info["userId"] # ............................................... 'p_created_by' + ), + session_token = session_token + ) + + # Done here: + return mongo_json["_id"] if mongo_json and db_json.get("status") == 1 else None + + async def set_token( + self, + db_conn: AsyncMySQL, + mongo_conn: AsyncMongo, + token_id: ObjectId | str, + client_user_id: dict, + token: dict, + session_token: str = None + ) -> bool: + + """ + This method is to be called when the end user authorizes your service to connect to his third-party account. For + example, when the end user allows you to access his GMail account. USE THIS FOR UPDATING (REFRESHING) TOKENS + ALSO. + :param db_conn: The database connection (MariaDB) to use to perform the action. + :param mongo_conn: The database connection (MongoDB) to use to perform the action. + :param token_id: The identifier granted by the 'get_token_id' method. + :param client_user_id: The way the third-party client recognizes your user. These details should match the + details furnished while requesting the authorization through 'get_token_id' method. + :param token: The token granted by the third-party service. + :param session_token: The session token of the user who requested this service. + :return: True if saved, False if failed. + """ + + # Start by assuming failure: + token_saved = False + + # Note down the timestamp at which this event occurred: + request_ts = date_time.get_current_utc_date_time(as_string = False) + + # Save the token to MongoDB: + mongo_json = await mongo_conn.find_one_and_update( + collection = MailOAuthModel.AUTH_COLLECTION, + filter = mongo_conn.dict_to_dot_notation({ + "_id": ObjectId(token_id), + "clientUserId": client_user_id + }), + update = [{ + "$set": { + "token": token, + "status": "active", + "lastRefreshTs": request_ts, + "firstRefreshTs": { + "$cond": { + "if": { + "$or": [ + {"$eq": ["$firstRefreshTs", None]}, + {"$eq": [{"$type": "$firstRefreshTs"}, "missing"]} + ] + }, + "then": request_ts, + "else": "$firstRefreshTs" + } + } + } + }], + projection = {"token": False}, + return_updated = True, + upsert = False + ) + + # Tell MariaDB that the token was saved: + if mongo_json is not None: + token_notes = { + "email": token["email"], + "displayName": token.get("displayName"), + "displayPictureUrl": token.get("displayPictureUrl"), + } + db_json = await self.call_procedure( + db_conn = db_conn, + proc_name = "entity_integration_save", + proc_args = ( + mongo_json["user"]["entityId"], # ............................... 'p_entity_id' + mongo_json["client"], # ......................................... 'p_provider' + "Active", # ..................................................... 'p_current_status' + "Auth Granted", # ............................................... 'p_last_action' + token["displayName"], # ......................................... 'p_display_name' + token["displayPictureUrl"], # ................................... 'p_display_picture' + token_id, # ..................................................... 'p_token_id' + json.to_string(python_data = token_notes, no_space = True), # ... 'p_notes' + mongo_json["user"]["userId"] # .................................. 'p_created_by' + ), + session_token = session_token + ) + if db_json["status"] == 1: token_saved = True + + # Done here: + return token_saved + + async def get_token( + self, + mongo_conn: AsyncMongo, + token_id: ObjectId | str = None, + **kwargs + ) -> dict | None: + + """ + To retrieve stored tokens from the database. + :param mongo_conn: The database connection (MongoDB) to use to perform the action. + :param token_id: The identifier granted by the 'get_token_id' method. + :param kwargs: Any set of key-value pairs to build custom search criteria. This could be things like the user + info, the client, the type of authentication used, or even the kind of service. + :return: The retrieved record that has the token, and information about the service and client if found, else + None when there is no matching record. + """ + + # Build the filter: + filter_json = {k: v for k, v in kwargs.items()} + if token_id: filter_json["_id"] = ObjectId(token_id) + + # If there is no search criteria, we exit with failure: + if not filter_json: return None + + # If there is some filtering possible, + # we fetch and return the token: + return await mongo_conn.find_one( + collection = self.AUTH_COLLECTION, + filter = filter_json, + projection = { + "_id": True, + "serviceType": True, + "authType": True, + "client": True, + "clientUserId": True, + "token": True + } + ) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/models/data/core/user_info.py b/models/data/core/user_info.py new file mode 100644 index 0000000..b391c20 --- /dev/null +++ b/models/data/core/user_info.py @@ -0,0 +1,212 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Saturday, 7th Dec., 2024. + + OBJECTIVE: + + To define how auth tokens will be stored in the database. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# For making data behaviour_models: +from pydantic import BaseModel, Field, field_validator, PastDatetime, AwareDatetime +from typing import Optional, Literal, Union + +# My utils: +from utils_v2.string import regex +from utils_v2.date_time import date_time + +# To work with MongoDB: +from bson.objectid import ObjectId + +# To work with date and time: +import datetime + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +class CoreAuthTokenModel(BaseModel): + + version: str = Field( + description = "a hint about the version no. of this message", + min_length = 1, + frozen = True, + default = "1.0.0" + ) + + serviceType: Literal["email", "sms", "chat"] = Field( + description = "the kind of service this message was sent/received from", + frozen = True + ) + + client: Literal[ + "gmail", "outlook", # ...................... Mail Clients + "telegram", "whatsapp", # .................. Chat Clients + "nimbusSmsIndia", "savvyBulkSmsKenya", # ... SMS Clients + "razorpay", "safaricomMPesaExpress" # ...... Payment Gateways + ] = Field( + description = "the third-part client that was used", + frozen = True + ) + + authType: Literal["oauth", "auth"] = Field( + description = "the type of authentication procedure used", + frozen = True + ) + + firstRequestTs: AwareDatetime = Field( + description = "the time (utc) at which authorization was first requested", + frozen = True + ) + + lastRequestTs: AwareDatetime = Field( + description = "the time (utc) at which authorization was last requested", + frozen = False + ) + + firstRefreshTs: AwareDatetime = Field( + description = "the time (utc) at which the tokens were first refreshed", + frozen = False + ) + + lastRefreshTs: AwareDatetime = Field( + description = "the time (utc) at which the tokens were last refreshed", + frozen = False + ) + + token: dict | None = Field( + description = "the actual auth tokens of that client; will differ for each client", + frozen = True, + default_factory = lambda: date_time.get_current_utc_date_time(as_string = False) + ) + + user: dict = Field( + description = "how you identify your user", + frozen = True + ) + + clientUserId: dict = Field( + description = "how third-party client identifies the same user", + frozen = True + ) + + status: Literal["active", "disabled"] = Field( + description = "to indicate the status of this account", + frozen = False, + default = "active" + ) + + syncFreq: Literal[60, 300, 1500] = Field( + description = "the no. of seconds after which to poll for updates from the client (if applicable)", + frozen = False, + default = 300 + ) + + # ┏┓ ┏• + # ┃ ┏┓┏┓╋┓┏┓ + # ┗┛┗┛┛┗┛┗┗┫ + # ┛ + + class Config: + extra = "allow" + arbitrary_types_allowed = True + + # ┓┏ ┓• ┓ • + # ┃┃┏┓┃┓┏┫┏┓╋┓┏┓┏┓ + # ┗┛┗┻┗┗┗┻┗┻┗┗┗┛┛┗ + + @field_validator( + "firstRequestTs", + "lastRequestTs", "firstRefreshTs", "lastRefreshTs", + mode = "before" + ) + def parse_date_time(cls, value): + return date_time.parse_date_time(input_value = value, timezone = date_time.TIMEZONE_UTC) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + from utils_v2.string import json + + auth_token = CoreAuthTokenModel( + serviceType = "email", + client = "gmail", + authType = "oauth", + firstRequestTs = date_time.get_current_utc_date_time(as_string = False), + lastRequestTs = date_time.get_current_utc_date_time(as_string = False), + firstRefreshTs = date_time.get_current_utc_date_time(as_string = False), + lastRefreshTs = date_time.get_current_utc_date_time(as_string = False), + token = { + "username": "testing123", + "password": "abcdefgh" + }, + user = { + "userId": 0, + "entityId": 1, + "billingAccountId": 2, + "fullName": "Bhopli" + }, + clientUserId = { + "email": "bhopli@gmail.com" + } + ) + + print("AUTH-TOKEN MODEL:", json.to_string(auth_token.model_dump(), default = str)) diff --git a/utils_v2/api/async_quart.py b/utils_v2/api/async_quart.py index 99915da..f5e5711 100644 --- a/utils_v2/api/async_quart.py +++ b/utils_v2/api/async_quart.py @@ -26,7 +26,7 @@ DECORATORS. """ - +import copy # ***************************************************************************************************************** # ***** **** @@ -81,6 +81,9 @@ import asyncio import time import datetime +# To make copies: +import copy + # ***************************************************************************************************************** # ***** **** @@ -374,55 +377,58 @@ def summarize_variable( :return: The summarized version of the input. """ - # If a null value was sent: + # Basic prep: if value is None: return None - - # Check the sensitive keys: if sensitive_keys is None: sensitive_keys = [] + value_copy = copy.deepcopy(value) + + # If the value is a Pydantic model, we convert to dict: + if isinstance(value_copy, pydantic.BaseModel): + value_copy = value_copy.model_dump() # Handle datatypes that you don't want to modify: - if isinstance(value, (int, float, bool, NoneType)): pass + if isinstance(value_copy, (int, float, bool, NoneType)): pass # When the value is a list or similar iterable: - elif isinstance(value, (list, tuple, set)): + elif isinstance(value_copy, (list, tuple, set)): if expand: if not isinstance(expand, bool): expand -= 1 - value = [AsyncLoggerContext.summarize( + value_copy = [AsyncLoggerContext.summarize( v, expand = expand, sensitive_keys = sensitive_keys - ) for v in value] - else: value = f"array of {len(value)} item(s)" + ) for v in value_copy] + else: value_copy = f"array of {len(value_copy)} item(s)" # If the value is a dict: - elif isinstance(value, dict): + elif isinstance(value_copy, dict): if expand: if not isinstance(expand, bool): expand -= 1 - value = { + value_copy = { k: AsyncLoggerContext.summarize( v, expand = expand, sensitive_keys = sensitive_keys - ) if k not in sensitive_keys else "********" - for k, v in value.items() + ) if k not in sensitive_keys else f"{len(str(v))} sensitive char(s)" + for k, v in value_copy.items() } - else: value = f"object of {len(value.keys())} field(s) [{', '.join(value.keys())}]" + else: value_copy = f"object of {len(value_copy.keys())} field(s) [{', '.join(value_copy.keys())}]" # When a dataframe is passed: - elif isinstance(value, pd.DataFrame): - cols = value.columns.to_list() - value = f"table with {len(cols)} col(s) [{', '.join(cols)}] and {len(value)} row(s)" - str_limit = 999 + elif isinstance(value_copy, pd.DataFrame): + cols = value_copy.columns.to_list() + value_copy = f"table with {len(cols)} col(s) [{', '.join(cols)}] and {len(value_copy)} row(s)" + str_limit = 1024 # If the input is some form of non-standard object: - else: value = str(value) + else: value_copy = str(value_copy) # Handle strings: - if isinstance(value, str): - if len(value) > str_limit: value = value[:str_limit] + "..." + if isinstance(value_copy, str): + if len(value_copy) > str_limit: value_copy = value_copy[:str_limit] + "...(trunc'd)" # Done here: - return value + return value_copy # --------------------------------------------------------------------------------------------------------------------- @@ -867,7 +873,7 @@ def log_request_to_mongo( if sensitive_keys: for k in sensitive_keys: for var in ["inbound_headers", "inbound_data"]: - try: kwargs[var][k] = len(str(kwargs[var][k])) * "*" + try: kwargs[var][k] = f"{len(str(kwargs[var][k]))} sensitive char(s)" except: pass # Construct the log: @@ -883,7 +889,7 @@ def log_request_to_mongo( ts = request_ts, tat = time.perf_counter() - start_ts, cpuTime = time.process_time() - cpu_start_ts, - sessionInfo = kwargs.get("session_info"), + sessionInfo = summarize_variable(kwargs.get("session_info"), expand = True), method = request_method, url = request_url, route = request_route,