(20241209) Better session info capture in the logs.

This commit is contained in:
2024-12-09 12:04:38 +05:30
parent 4583a5e63c
commit 30088dcea4
5 changed files with 792 additions and 24 deletions
View File
+221
View File
@@ -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())
+329
View File
@@ -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
+212
View File
@@ -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))