diff --git a/api/blueprints/sms/auth.py b/api/blueprints/sms/auth_v2.py similarity index 100% rename from api/blueprints/sms/auth.py rename to api/blueprints/sms/auth_v2.py diff --git a/api/blueprints/sms/list.py b/api/blueprints/sms/list.py new file mode 100644 index 0000000..6ae47cf --- /dev/null +++ b/api/blueprints/sms/list.py @@ -0,0 +1,242 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 19th Dec., 2024 + + OBJECTIVE: + + To receive auth details for various SMS client APIs. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + + NOTES: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# For using Quart: +from quart import Blueprint, current_app, request + +# My utils: +from utils_v2.string import json +from utils_v2.api.codes import StatusCodes, HttpCodes +from utils_v2.api.response import ResponseModel +from utils_v2.api.async_quart import ( + set_api_version, + read_input, + get_session_info, + log_request_to_mongo, + log_chain_to_mongo, + should_not_be_under_maintenance, + only_whitelisted_ips, + limit_rate, + validate_input, + handle_cancelled_request +) + +# Common: +from shared import constants + +# Data Models: +from models.api.sms.auth import SMSAuthRequestHeaders, SMSAuthRequestData +from models.core.auth_token import CoreAuthTokenModel + +# For asynchronous activities: +import asyncio + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# Related to Quart: +sms_auth_bp = Blueprint("sms_auth", __name__) + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +@sms_auth_bp.record_once +def init(blueprint_setup_state): + + # This gets called when the blueprint is registered. + # Consider this to be a one-time setup for the whole blueprint: + pass + + +# --------------------------------------------------------------------------------------------------------------------- + + +@sms_auth_bp.route("/auth", methods = ["POST"]) +@set_api_version(api_version = "1.0.0") +@read_input(sanitize_headers = False, sanitize_data = False) +@get_session_info(key = "X-Session-Token", session_coro = "get_session") +@log_request_to_mongo( + attr_name = "logs_mongo", + project = constants.PROJECT_NAME, + log_type = constants.MODULE_NAME, + operation = "smsAuthApi", + log_input = True, + log_output = True, + sensitive_keys = ["sessionToken", "X-Session-Token"] +) +@log_chain_to_mongo(attr_name = "logs_mongo") +@should_not_be_under_maintenance(attr_name = "is_under_maintenance") +@validate_input( + header_validator = lambda x: SMSAuthRequestHeaders(**x).model_dump(), + data_validator = lambda x: SMSAuthRequestData(**x) +) +@handle_cancelled_request() +async def authorize_sms_client( + inbound_headers: dict | SMSAuthRequestHeaders = None, + inbound_data: dict | SMSAuthRequestData = None, + inbound_files: dict = None, + **kwargs +): + + """ + Use this when a user wants to register a third-party SMS client with your service. + :param inbound_headers: auto-extracted by the decorators. + :param inbound_data: auto-extracted by the decorators. + :param inbound_files: auto-extracted by the decorators. + :param kwargs: Any number of extra inputs supplied by the decorators. + :return: A standard response structure. + """ + + # ┏┓ + # ┃┃┏┓┏┓┏┓┏┓┏┓┏┏┓┏┏ + # ┣┛┛ ┗ ┣┛┛ ┗┛┗┗ ┛┛ + # ┛ + + # If the session token is invalid/expired: + if kwargs.get("session_info") is None: + return ResponseModel( + status_code = StatusCodes.FAILED, + http_code = HttpCodes.UNAUTHORIZED + ) + + # Start by assuming failure: + success = False + + # ┏┓ ┳┓• ┓ ┏┓┳┳┓┏┓ ┳ ┓• + # ┣ ┏┓┏┓ ┃┃┓┏┳┓┣┓┓┏┏ ┗┓┃┃┃┗┓ ┃┏┓┏┫┓┏┓ + # ┻ ┗┛┛ ┛┗┗┛┗┗┗┛┗┻┛ ┗┛┛ ┗┗┛ ┻┛┗┗┻┗┗┻ + + if inbound_data.smsClient == "nimbusSmsIndia": + + success = await current_app.sms_controller.set_token_direct( + sql_conn = current_app.sql_writer, + mongo_data_conn = current_app.data_mongo, + auth_token = CoreAuthTokenModel( + serviceType = "sms", + client = inbound_data.smsClient, + authType = "auth", + auth = inbound_data.auth.model_dump(), + user = kwargs.get("session_info"), + clientUserId = { + "userId": inbound_data.auth.userId, + "senderId": inbound_data.auth.senderId, + "entityId": inbound_data.auth.entityId + }, + status = "active", + syncFreq = 60 + ), + token_notes = {}, + session_token = inbound_headers["X-Session-Token"] + ) + + # ┏┓ ┏┓ ┳┓ ┓┓ ┏┓┳┳┓┏┓ ┓┏┓ + # ┣ ┏┓┏┓ ┗┓┏┓┓┏┓┏┓┏ ┣┫┓┏┃┃┏ ┗┓┃┃┃┗┓ ┃┫ ┏┓┏┓┓┏┏┓ + # ┻ ┗┛┛ ┗┛┗┻┗┛┗┛┗┫ ┻┛┗┻┗┛┗ ┗┛┛ ┗┗┛ ┛┗┛┗ ┛┗┗┫┗┻ + # ┛ ┛ + + elif inbound_data.smsClient == "savvyBulkSmsKenya": + + success = await current_app.sms_controller.set_token_direct( + sql_conn = current_app.sql_writer, + mongo_data_conn = current_app.data_mongo, + auth_token = CoreAuthTokenModel( + serviceType = "sms", + client = inbound_data.smsClient, + authType = "auth", + auth = inbound_data.auth.model_dump(), + user = kwargs.get("session_info"), + clientUserId = { + "partnerId": inbound_data.auth.partnerId, + "shortCode": inbound_data.auth.shortCode + }, + status = "active", + syncFreq = 60 + ), + token_notes = {}, + session_token = inbound_headers["X-Session-Token"] + ) + + # ┳┓ + # ┣┫┏┓┏┏┓┏┓┏┓┏┏┓ + # ┛┗┗ ┛┣┛┗┛┛┗┛┗ + # ┛ + + # Done here: + return ResponseModel( + status_code = StatusCodes.OK if success else StatusCodes.FAILED, + http_code = HttpCodes.SUCCESS if success else HttpCodes.INTERNAL_SERVER_ERROR, + data = { + "client": inbound_data.smsClient, + "authorized": success + } + ) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/api/blueprints/sms/send.py b/api/blueprints/sms/send_v2.py similarity index 100% rename from api/blueprints/sms/send.py rename to api/blueprints/sms/send_v2.py diff --git a/api/blueprints/sms/tags.py b/api/blueprints/sms/tags.py new file mode 100644 index 0000000..ca81bc2 --- /dev/null +++ b/api/blueprints/sms/tags.py @@ -0,0 +1,230 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 19th Dec., 2024 + + OBJECTIVE: + + To list SMS messages associated with incoming identifiers. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + + NOTES: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys + +from bson import ObjectId + +sys.path.append(".") +sys.path.append("..") + +# For using Quart: +from quart import Blueprint, current_app, request + +# My utils: +from utils_v2.string import json +from utils_v2.api.codes import StatusCodes, HttpCodes +from utils_v2.api.response import ResponseModel +from utils_v2.api.async_quart import ( + set_api_version, + read_input, + get_session_info, + log_request_to_mongo, + log_chain_to_mongo, + should_not_be_under_maintenance, + only_whitelisted_ips, + limit_rate, + validate_input, + handle_cancelled_request +) + +# Common: +from shared import constants + +# Data Models: +from models.core.user import CoreUserInfoModel +from models.api.sms.list import SMSListRequestHeaders, SMSListRequestData + +# Helpers: +from api.helpers.user import token_check + +# For asynchronous activities: +import asyncio + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# Related to Quart: +sms_list_bp = Blueprint("sms_list", __name__) + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +@sms_list_bp.record_once +def init(blueprint_setup_state): + + # This gets called when the blueprint is registered. + # Consider this to be a one-time setup for the whole blueprint: + pass + + +# --------------------------------------------------------------------------------------------------------------------- + + +@sms_list_bp.route("/list", methods = ["GET"]) +@set_api_version(api_version = "1.0.0") +@read_input(sanitize_headers = False, sanitize_data = False) +@get_session_info(key = "X-Session-Token", session_coro = "get_session") +@log_request_to_mongo( + attr_name = "logs_mongo", + project = constants.PROJECT_NAME, + log_type = constants.MODULE_NAME, + operation = "smsListApi", + log_input = True, + log_output = True, + sensitive_keys = ["sessionToken", "X-Session-Token", "tokenKeys"] +) +@log_chain_to_mongo(attr_name = "logs_mongo") +@should_not_be_under_maintenance(attr_name = "is_under_maintenance") +@validate_input( + header_validator = lambda x: SMSListRequestHeaders(**x).model_dump(), + data_validator = lambda x: SMSListRequestData(**x) +) +@handle_cancelled_request() +async def list_sms_messages( + inbound_headers: dict | SMSListRequestHeaders = None, + inbound_data: dict | SMSListRequestData = None, + inbound_files: dict = None, + **kwargs +): + + """ + Use this API when a user wants his SMS messages listed on the screen. + :param inbound_headers: auto-extracted by the decorators. + :param inbound_data: auto-extracted by the decorators. + :param inbound_files: auto-extracted by the decorators. + :param kwargs: Any number of extra inputs supplied by the decorators. + :return: A standard response structure. + """ + + # ┏┓ ┓ ┏┓┓ ┓ + # ┣┫┓┏╋┣┓ ┃ ┣┓┏┓┏┃┏ + # ┛┗┗┻┗┛┗ ┗┛┛┗┗ ┗┛┗ + + # If the session token is invalid/expired: + if kwargs.get("session_info") is None: + return ResponseModel( + status_code = StatusCodes.FAILED, + http_code = HttpCodes.UNAUTHORIZED + ) + + # ┏┓ ┓ • ┏┓┓ ┓ + # ┃┃┓┏┏┏┓┏┓┏┓┏┣┓┓┏┓ ┃ ┣┓┏┓┏┃┏ + # ┗┛┗┻┛┛┗┗ ┛ ┛┛┗┗┣┛ ┗┛┛┗┗ ┗┛┗ + # ┛ + + # Get the tokens from the database: + auth_tokens = await current_app.sms_controller.get_tokens_from_keys( + mongo_data_conn = current_app.data_mongo, + token_keys = inbound_data.tokenKeys, + limit = len(inbound_data.tokenKeys) + ) + token_ids = [ObjectId(t.authTokenId) for t in auth_tokens] + + # Check if these tokens belong to the user claiming ownership: + if not await token_check.is_authorized( + mongo_conn = current_app.data_mongo, + user_info = CoreUserInfoModel(**kwargs["session_info"]), + token_ids = token_ids + ): return ResponseModel( + status_code = StatusCodes.FAILED, + http_code = HttpCodes.UNAUTHORIZED, + message = "User doesn't have rights over one or more SMS accounts." + ) + + # ┳┓ ┓ • • + # ┃┃┏┓╋┏┓ ┃ ┓┏╋┓┏┓┏┓ + # ┻┛┗┻┗┗┻ ┗┛┗┛┗┗┛┗┗┫ + # ┛ + + messages = await current_app.sms_controller.get_messages( + mongo_data_conn = current_app.data_mongo, + token_ids = token_ids, + limit = inbound_data.count, + skip = inbound_data.fromCount, + projection = { + "message.metadata": False, + "message.rawResponse": False + } + ) + + # ┳┓ + # ┣┫┏┓┏┏┓┏┓┏┓┏┏┓ + # ┛┗┗ ┛┣┛┗┛┛┗┛┗ + # ┛ + + # Done here: + message_count = len(messages) + success = True if messages is not None and message_count > 0 else False + return ResponseModel( + status_code = StatusCodes.OK if success else StatusCodes.FAILED, + http_code = HttpCodes.SUCCESS if success else HttpCodes.NOT_FOUND, + data = [m.full for m in messages] if success else None, + message = f"{message_count} SMS message(s) found." + ) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/utils_v2/goog/gmail/__init__.py b/controllers_v2/__init__.py similarity index 100% rename from utils_v2/goog/gmail/__init__.py rename to controllers_v2/__init__.py diff --git a/controllers_v2/core/__init__.py b/controllers_v2/core/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/controllers_v2/core/ai/__init__.py b/controllers_v2/core/ai/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/controllers_v2/core/ai/llm.py b/controllers_v2/core/ai/llm.py new file mode 100644 index 0000000..92c543b --- /dev/null +++ b/controllers_v2/core/ai/llm.py @@ -0,0 +1,180 @@ +""" + + 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.database.async_mongo_v2 import AsyncMongo + +# Base model: +from controllers.base import BaseModel + +# Data Models: +from models.core.user import CoreUserInfoModel +from models.core.ai.llm import LLMInput, LLMOutput, LLMUsageTokens + +# To work with LLMs: +from langchain_openai import ChatOpenAI + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class CoreLLMController(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: CoreUserInfoModel, + 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.model_dump()} + 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 + ) + if inserted_id: llm_response.invocationId = str(inserted_id) + + # Done here: + return llm_response + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/controllers_v2/core/auth_token.py b/controllers_v2/core/auth_token.py new file mode 100644 index 0000000..0bf52cd --- /dev/null +++ b/controllers_v2/core/auth_token.py @@ -0,0 +1,381 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 12th Dec., 2024 + + OBJECTIVE: + + To handle all auth-tokens from one place. + + 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 controllers.base import BaseModel + +# Data models: +from models.core.auth_token import CoreAuthTokenModel + +# To work with datatypes: +from typing import List + +# To work with MongoDB: +from bson import ObjectId + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class CoreAuthTokenController(BaseModel): + + # ┏┓┓ ┓┏ + # ┃ ┃┏┓┏┏ ┃┃┏┓┏┓┏ + # ┗┛┗┗┻┛┛ ┗┛┗┻┛ ┛ + + # For MongoDB: + AUTH_COLLECTION = "_authTokens" + + async def get_token_key( + self, + db_conn: AsyncMySQL, + mongo_conn: AsyncMongo, + auth_token: CoreAuthTokenModel, + token_notes: dict, + 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 auth_token: An instance of the core auth-token model that holds data in the database. + :param token_notes: Any notes to feed into MariaDB with the token identifier. + :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 if it already exists, else create one. + # BE CAREFUL WITH THE KEYS HERE, THEY SHOULD MATCH THE FIELDS OF THE CORE AUTH-TOKEN MODEL: + mongo_json = await mongo_conn.find_one_and_update( + collection = self.AUTH_COLLECTION, + filter = mongo_conn.dict_to_dot_notation({ + "serviceType": auth_token.serviceType, + "user": { + "entityId": auth_token.user.entityId, + "billingAccountId": auth_token.user.billingAccountId + }, + "clientUserId": auth_token.clientUserId + }), + update = { + "$set": { + "lastRequestTs": auth_token.lastRequestTs, + "status": auth_token.status, + "syncFreq": auth_token.syncFreq + }, + "$setOnInsert": { + "key": auth_token.key, + "serviceType": auth_token.serviceType, + "client": auth_token.client, + "authType": auth_token.authType, + "user": auth_token.user.model_dump(), + "clientUserId": auth_token.clientUserId, + "auth": auth_token.auth, + "token": auth_token.token, + "firstRefreshTs": auth_token.firstRefreshTs, + "lastRefreshTs": auth_token.lastRefreshTs, + "firstRequestTs": auth_token.firstRequestTs or request_ts, + } + }, + projection = { + "_id": True, + "key": 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 = ( + auth_token.user.entityId, # ..................................... 'p_entity_id' + auth_token.client, # ............................................ 'p_provider' + auth_token.status, # ............................................ 'p_current_status' + "Auth Requested", # ............................................. 'p_last_action' + None, # ......................................................... 'p_display_name' + None, # ......................................................... 'p_display_picture' + mongo_json["key"], # ............................................ 'p_token_id' + json.to_string(python_data = token_notes, no_space = True), # ... 'p_notes' + auth_token.user.userId # ........................................ 'p_created_by' + ), + session_token = session_token + ) + + # Done here: + return mongo_json["key"] if mongo_json and db_json.get("status") == 1 else None + + async def set_token( + self, + db_conn: AsyncMySQL, + mongo_conn: AsyncMongo, + token_key: ObjectId | str, + auth_token: CoreAuthTokenModel, + token_notes: 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_key: The identifier granted by the 'get_token_key' method. + :param auth_token: The actual auth/token data to be saved to the database. + :param token_notes: Any notes to feed into MariaDB with the token identifier. + :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. + # BE CAREFUL WITH THE KEYS HERE, THEY SHOULD MATCH THE FIELDS OF THE CORE AUTH-TOKEN MODEL: + mongo_json = await mongo_conn.find_one_and_update( + collection = self.AUTH_COLLECTION, + filter = mongo_conn.dict_to_dot_notation({ + "key": ObjectId(token_key), + "clientUserId": auth_token.clientUserId + }), + update = [{ + "$set": { + "auth": auth_token.auth, + "token": auth_token.token, + "status": auth_token.status, + "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: + 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' + auth_token.status, # ............................................ 'p_current_status' + "Auth Granted", # ............................................... 'p_last_action' + auth_token.token.get("displayName"), # .......................... 'p_display_name' + auth_token.token.get("displayPictureUrl"), # .................... 'p_display_picture' + token_key, # .................................................... 'p_token_id' + json.to_string(python_data = token_notes, no_space = True), # ... 'p_notes' + auth_token.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_from_id( + self, + mongo_conn: AsyncMongo, + token_id: ObjectId | str = None, + ) -> CoreAuthTokenModel | None: + + """ + To retrieve stored tokens from the database. One token at a time. + :param mongo_conn: The database connection (MongoDB) to use to perform the action. + :param token_id: The identifier of the document that holds the token's details. + :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. + """ + + # If there is some filtering possible, we fetch the token: + token = await mongo_conn.find_one( + collection = self.AUTH_COLLECTION, + filter = {"_id": ObjectId(token_id)} + ) + + # Done here: + return CoreAuthTokenModel(**token) if token else None + + async def get_token_from_key( + self, + mongo_conn: AsyncMongo, + token_key: ObjectId | str = None, + ) -> CoreAuthTokenModel | None: + + """ + To retrieve stored tokens from the database. One token at a time. + :param mongo_conn: The database connection (MongoDB) to use to perform the action. + :param token_key: The identifier granted by the 'get_token_key' method. + :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. + """ + + # If there is some filtering possible, we fetch the token: + token = await mongo_conn.find_one( + collection = self.AUTH_COLLECTION, + filter = {"key": ObjectId(token_key)} + ) + + # Done here: + return CoreAuthTokenModel(**token) if token else None + + async def get_tokens_from_ids( + self, + mongo_conn: AsyncMongo, + token_ids: List[ObjectId | str] = None, + limit: int = 100 + ) -> List[CoreAuthTokenModel]: + + """ + To retrieve stored tokens from the database. Multiple tokens at a time. + :param mongo_conn: The database connection (MongoDB) to use to perform the action. + :param token_ids: The identifier of the document that holds the token's details. + :param limit: The max. no. of records to pick. + :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. + """ + + # If there is some filtering possible, we fetch the token: + tokens = await mongo_conn.find_many( + collection = self.AUTH_COLLECTION, + filter = {"_id": {"$in": [ObjectId(k) for k in token_ids]}}, + limit = limit + ) + + # Done here: + return [CoreAuthTokenModel(**token) for token in tokens] + + async def get_tokens_from_keys( + self, + mongo_conn: AsyncMongo, + token_keys: List[ObjectId | str] = None, + limit: int = 100 + ) -> List[CoreAuthTokenModel]: + + """ + To retrieve stored tokens from the database. Multiple tokens at a time. + :param mongo_conn: The database connection (MongoDB) to use to perform the action. + :param token_keys: the identifiers granted by the 'get_token_key' method. + :param limit: The max. no. of records to pick. + :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. + """ + + # If there is some filtering possible, we fetch the token: + tokens = await mongo_conn.find_many( + collection = self.AUTH_COLLECTION, + filter = {"key": {"$in": [ObjectId(k) for k in token_keys]}}, + limit = limit + ) + + # Done here: + return [CoreAuthTokenModel(**token) for token in tokens] + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/controllers_v2/core/base.py b/controllers_v2/core/base.py new file mode 100644 index 0000000..6ba2ec3 --- /dev/null +++ b/controllers_v2/core/base.py @@ -0,0 +1,388 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Tuesday, 22nd Oct., 2024 + + OBJECTIVE: + + To provide an easy way to create models to handle documents for Bicree. + This is the base model for this microservice. It will define the structure for all other models that will be + used in this particular microservice. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# My utils: +from utils_v2.database.async_mysql_v2 import AsyncMySQL +from utils_v2.cache.async_redis_cache_v2 import AsyncRedisCache +from utils_v2.api.codes import StatusCodes + +# For asynchronous activities: +import asyncio + +# For debugging: +from icecream import IceCreamDebugger + +# To work with datatypes: +from typing import List + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class BaseModel: + + PREVIEW_LENGTH = 250 + + def __init__( + self, + cache = None, + alert_url = None, + http_client = None, + debug = True, + debug_prefix = "Model | ", + debug_only_errors = True + ): + + """ + This is the base model. + :param cache: The object to use for caching results from database calls. + :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. + """ + + # Prepare the caching utility: + self._cache = cache + + # For sending alerts: + self._alert_url = alert_url + self._http_client = http_client + + # Prepare the debugging utility: + self._debug_prefix = debug_prefix + self._printer = IceCreamDebugger(prefix = debug_prefix, includeContext = True) + if not debug: self._printer.disable() + self._debug_only_errors = debug_only_errors + + # A semaphore for activities that must absolutely be done one at a time: + self.__exclusive_semaphore = asyncio.Semaphore(1) + + # A simple debugging output: + self._printer("Model initialized.") + + def enable_terminal_print(self): + self._printer.enable() + + def disable_terminal_print(self): + self._printer.disable() + + def debug_only_errors(self): + self._debug_only_errors = True + + def debug_everything(self): + self._debug_only_errors = False + + async def send_alert( + self, + message: str, + session_token = None, + alert_type = "error" + ): + + """ + Sends out an alert (ideally through the tech module). This is meant to be used when some exception occurs, and + you want to be informed before the client complains. + :param message: The message to send out to the admins. + :param session_token: The session token of the user (optional) so that the alert message can display the name of + the user who faced the trouble. + :param alert_type: The type of alert to throw ("error", "warning", or "info"). + :return: None. + """ + + if self._http_client is not None and self._alert_url is not None: + response = await self._http_client.post( + url = self._alert_url, + headers = {"X-Session-Token": session_token} if session_token else None, + json = { + "message": message, + "type": alert_type + } + ) + + async def call_cached_procedure( + self, + cache: AsyncRedisCache, + cache_key: str, + cache_expiry: int, + db_conn: AsyncMySQL, + proc_name: str, + proc_args: tuple, + retry_count: int = 1, + backoff_seconds: float = 0.5, + backoff_multiplier: float = 1.1, + session_token: str = None + ): + + """ + Calls a stored procedure and returns the response as a JSON-like object (dict or list). + :param cache: The caching object to use to set the session in cache memory. + :param cache_key: The string to use as the key when caching the response. + :param cache_expiry: The no. of seconds after which this information will be deleted from the cache. + :param db_conn: The connection instance to use to call the procedure. + :param proc_name: The name of the stored procedure that must be called. + :param proc_args: The args to be sent to the stored procedure. + :param retry_count: The max. number of times to try in case one or more attempts fail. + :param backoff_seconds: The time to wait before making the next attempt if the retry count is more than 1. + :param backoff_multiplier: The factor that dictates how much to modify the time delay by when waiting to retry. + :param session_token: A session token to share with the tech module when alerts need to be sent out for any + occurrence of exceptions. If this is passed, the tech module will be able to tell you which user faced the + issue. + :return: The response from the stored procedure. + """ + + # check for the data in cache: + data = await cache.get(cache_key) + + # If the data isn't in the cache, call the procedure: + if data is None: + + # Make the database call: + data = await self.call_procedure( + db_conn = db_conn, + proc_name = proc_name, + proc_args = proc_args, + retry_count = retry_count, + backoff_seconds = backoff_seconds, + backoff_multiplier = backoff_multiplier, + session_token = session_token + ) + + # If the database call succeeded, cache the response: + if isinstance(data, dict) and data["status"] == 1: + await cache.set(key = cache_key, value = data, expiry = cache_expiry) + + # Done here: + return data + + async def call_procedure( + self, + db_conn: AsyncMySQL, + proc_name: str, + proc_args: tuple, + retry_count: int = 1, + backoff_seconds: float = 0.5, + backoff_multiplier: float = 1.1, + session_token: str = None + ): + + """ + Calls a stored procedure and returns the response as a JSON-like object (dict or list). + :param db_conn: The connection instance to use to call the procedure. + :param proc_name: The name of the stored procedure that must be called. + :param proc_args: The args to be sent to the stored procedure. + :param retry_count: The max. number of times to try in case one or more attempts fail. + :param backoff_seconds: The time to wait before making the next attempt if the retry count is more than 1. + :param backoff_multiplier: The factor that dictates how much to modify the time delay by when waiting to retry. + :param session_token: A session token to share with the tech module when alerts need to be sent out for any + occurrence of exceptions. If this is passed, the tech module will be able to tell you which user faced the + issue. + :return: The response from the stored procedure. + """ + + # Call the stored procedure: + db_json, exception = await db_conn.call_procedure_and_get_json( + proc_name, + proc_args, + retry_count = retry_count, + backoff_seconds = backoff_seconds, + backoff_multiplier = backoff_multiplier, + return_exception = True + ) + + # Understand the response: + success = True if db_json["status"] == 1 else False + message = db_json.get("message") + + # Debugging print: + if not success or not self._debug_only_errors: + self._printer(proc_name, proc_args, success, exception, message) + + # Send an alert out on exceptions: + if exception is not None: + + # Format the message in Markdown format: + exception_string = str(exception).replace("`", "'") + formatted_message = f"*Module:*\n`{self._debug_prefix}`\n\n" + formatted_message += f"*Proc:*\n`{proc_name}`\n\n" + formatted_message += f"*Args:*\n`({', '.join([str(_) for _ in proc_args])})`\n\n" + formatted_message += f"*Arg-Types:*\n`({', '.join([type(_).__name__ for _ in proc_args])})`\n\n" + formatted_message += f"*Message:*\n`{message}`\n\n" + formatted_message += f"*Success:*\n`{success}`\n\n" + formatted_message += f"*Exception:*\n`{exception_string}`\n\n" + + # Send the alert: + await self.send_alert(formatted_message, session_token = session_token) + + # Return the response: + db_json["status_code"] = StatusCodes.OK if success else StatusCodes.FAILED + return db_json + + async def execute_one( + self, + db_conn: AsyncMySQL, + query: str, + session_token: str = None + ): + + """ + Runs one query and sends an alert if that fails. + :param db_conn: The connection to use to run the query. + :param query: The query to run. + :param session_token: A session token to share with the tech module when alerts need to be sent out for any + occurrence of exceptions. If this is passed, the tech module will be able to tell you which user faced the + issue. + :return: The response from the database. + """ + + # Run the query: + rows_affected, db_response, exception = await db_conn.execute_one(query = query, return_exception = True) + + # Send an alert out on exceptions: + if exception is not None: + + # Created needed previews: + query_preview = query if len(query) <= self.PREVIEW_LENGTH else query[:self.PREVIEW_LENGTH] + "..." + + # Format the message in Markdown format: + formatted_message = f"*Module:*\n`{self._debug_prefix}`\n\n" + formatted_message += f"*Query:*\n`{query_preview}`\n\n" + formatted_message += f"*Rows Affected:*\n`{rows_affected}`\n\n" + formatted_message += f"*DB Response:*\n`{db_response}`\n\n" + formatted_message += f"*Exception:*\n`{exception}`\n\n" + + # Send the alert: + await self.send_alert(formatted_message, session_token = session_token) + + # Return the response: + return rows_affected, db_response + + async def execute_many( + self, + db_conn: AsyncMySQL, + query: str, + data: List[tuple], + session_token: str = None + ): + + """ + Runs many queries and sends an alert if that fails. + :param db_conn: The connection to use to run the query. + :param query: The query to run. + :param data: The data to feed into the query. + :param session_token: A session token to share with the tech module when alerts need to be sent out for any + occurrence of exceptions. If this is passed, the tech module will be able to tell you which user faced the + issue. + :return: The response from the database. + """ + + # Run the query: + rows_affected, db_response, exception = await db_conn.execute_many( + query = query, + data = data, + return_exception = True + ) + + # Send an alert out on exceptions: + if exception is not None: + + # Created needed previews: + query_preview = query if len(query) <= self.PREVIEW_LENGTH else query[:self.PREVIEW_LENGTH] + "..." + data_preview = str(data) + if len(data_preview) > self.PREVIEW_LENGTH: data_preview = data_preview[:self.PREVIEW_LENGTH] + "..." + + # Format the message in Markdown format: + formatted_message = f"*Module:*\n`{self._debug_prefix}`\n\n" + formatted_message += f"*Query:*\n`{query_preview}`\n\n" + formatted_message += f"*Data:*\n`{data_preview}`\n\n" + formatted_message += f"*Rows Affected:*\n`{rows_affected}`\n\n" + formatted_message += f"*DB Response:*\n`{db_response}`\n\n" + formatted_message += f"*Exception:*\n`{exception}`\n\n" + + # Send the alert: + await self.send_alert(formatted_message, session_token = session_token) + + # Return the response: + return rows_affected, db_response + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/controllers_v2/core/message.py b/controllers_v2/core/message.py new file mode 100644 index 0000000..257d7bb --- /dev/null +++ b/controllers_v2/core/message.py @@ -0,0 +1,399 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 12th Dec., 2024 + + OBJECTIVE: + + To handle all messages from one place. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# For Quart: +from quart import current_app + +# 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, AsyncMongoStorage + +# Base model: +from controllers.base import BaseModel + +# Data models: +from models.core.auth_token import CoreAuthTokenModel +from models.core.message import CoreMessageModel +from models.core.user import CoreUserInfoModel + +# To work with MongoDB: +from bson import ObjectId +from pymongo import InsertOne, UpdateOne, ReplaceOne + +# To work with datatypes: +from typing import Literal, List, Dict, Any + +# To make deep-copies: +import copy + +# To work with base-64 encoding: +import base64 + +# To work with date and time: +import datetime + +# For asynchronous activities: +import asyncio + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class CoreMessageController(BaseModel): + + # ┏┓┓ ┓┏ + # ┃ ┃┏┓┏┏ ┃┃┏┓┏┓┏ + # ┗┛┗┗┻┛┛ ┗┛┗┻┛ ┛ + + # For MongoDB: + MESSAGES_COLLECTION = "_messages" + + # ┏┓┳┓┳┳┳┓ ┏┓ + # ┃ ┣┫┃┃┃┃ ━━ ┃ ┏┓┏┓┏┓╋┏┓ + # ┗┛┛┗┗┛┻┛ ┗┛┛ ┗ ┗┻┗┗ + + async def insert( + self, + mongo_conn: AsyncMongo, + message: CoreMessageModel + ) -> ObjectId: + + """ + Simply insert one message document into the database. + :param mongo_conn: The instance of the database connector to use for the operation. + :param message: The message to save into the database. + :return: The object id of the inserted document. + """ + + # Simply insert the document: + return await mongo_conn.insert_one( + collection = self.MESSAGES_COLLECTION, + document = message, + raise_exception = True + ) + + async def bulk_write( + self, + mongo_conn: AsyncMongo, + mongo_operations: list + ) -> int: + + """ + Needed in cases like forcing re-sync of mails where you need to perform actions like bulk replacements of + existing documents. Not recommended to use. Please use very carefully to ensure document integrity. + :param mongo_conn: The instance of the database connector to use for the operation. + :param mongo_operations: The list operations that are supported by MongoDB's Bulk Write system. + :return: The no. of documents affected. + """ + + return await mongo_conn.bulk_write( + collection = self.MESSAGES_COLLECTION, + requests = mongo_operations, + raise_exception = True + ) + + # ┏┓┳┓┳┳┳┓ ┳┓ • + # ┃ ┣┫┃┃┃┃ ━━ ┣┫┏┓╋┏┓┓┏┓┓┏┏┓ + # ┗┛┛┗┗┛┻┛ ┛┗┗ ┗┛ ┗┗ ┗┛┗ + + async def count_messages( + self, + mongo_conn: AsyncMongo, + token_ids: List[ObjectId | str], + additional_filter: dict = None + ) -> int: + + """ + Just counts the no. of messages that match a given set of conditions. + :param mongo_conn: The instance of the database connector to use for the operation. + :param token_ids: The token ids of the accounts from which these messages must be fetched. + :param additional_filter: Any addition filters to use. + :return: The no. of messages that match the given conditions. + """ + + # Prepare the filter: + if not isinstance(token_ids, list): token_ids = [token_ids] + token_ids = [ObjectId(t) for t in token_ids] + filter_json = {"tokenId": {"$in": token_ids}} + if additional_filter: + for k, v in additional_filter.items(): + filter_json[k] = v + + # Get the count of the documents that match the criteria: + count = await mongo_conn.count( + collection = self.MESSAGES_COLLECTION, + filter = filter_json, + raise_exception = True + ) + + # Done here: + return count + + async def get_previews( + self, + mongo_conn: AsyncMongo, + token_ids: List[ObjectId | str], + limit: int = 100, + skip: int = 0, + additional_filter: dict = None + ) -> List[CoreMessageModel] | None: + + """ + Fetches many messages in one call, but leaves out the full payloads. + :param mongo_conn: The instance of the database connector to use for the operation. + :param token_ids: The token ids of the accounts from which these messages must be fetched. + :param limit: The max. no. of messages to retrieve in this call. + :param skip: The no. of initial messages to skip. Useful for pagination. + :param additional_filter: Any addition filters to use. + :return: The list of messages (as the message model). This list can be empty. + """ + + # Prepare the filter: + if not isinstance(token_ids, list): token_ids = [token_ids] + token_ids = [ObjectId(t) for t in token_ids] + filter_json = {"tokenId": {"$in": token_ids}} + if additional_filter: + for k, v in additional_filter.items(): + filter_json[k] = v + + # We fetch the messages that are identified by a specific token id, + # with the specified fetching limits, while enforcing the sorting condition: + records = await mongo_conn.find_many( + collection = self.MESSAGES_COLLECTION, + filter = filter_json, + limit = limit, + skip = skip, + sort = {"ts": -1}, + projection = { + "_id": True, + "ts": True, + "syncTs": True, + "tokenId": True, + "serviceType": True, + "client": True, + "clientMessageId": True, + "clientThreadId": True, + "isSent": True, + "isBroadcast": True, + "sentSuccessfully": True, + "sender": True, + "chat": True, + "snippet": True, + "aiSnippet": True, + "tags": True + }, + raise_exception = True + ) + + # Convert the fetched records to instances of the data model and return: + for record in records: record["message"] = {} + return [CoreMessageModel(**record) for record in records] + + async def get_messages( + self, + mongo_conn: AsyncMongo, + token_ids: List[ObjectId | str], + limit: int = 100, + skip: int = 0, + additional_filter: dict = None + ) -> List[CoreMessageModel] | None: + + """ + Fetches many full messages in one call. + :param mongo_conn: The instance of the database connector to use for the operation. + :param token_ids: The token ids of the accounts from which these messages must be fetched. + :param limit: The max. no. of messages to retrieve in this call. + :param skip: The no. of initial messages to skip. Useful for pagination. + :param additional_filter: Any addition filters to use. + :return: The list of messages (as the message model). This list can be empty. + """ + + # Prepare the filter: + if not isinstance(token_ids, list): token_ids = [token_ids] + token_ids = [ObjectId(t) for t in token_ids] + filter_json = {"tokenId": {"$in": token_ids}} + if additional_filter: + for k, v in additional_filter.items(): + filter_json[k] = v + + # We fetch the messages that are identified by a specific token id, + # with the specified fetching limits, while enforcing the sorting condition: + records = await mongo_conn.find_many( + collection = self.MESSAGES_COLLECTION, + filter = filter_json, + limit = limit, + skip = skip, + sort = {"ts": -1}, + raise_exception = True + ) + + # Convert the fetched records to instances of the data model and return: + return [CoreMessageModel(**record) for record in records] + + async def get_message( + self, + mongo_conn: AsyncMongo, + message_id: ObjectId | str, + ) -> CoreMessageModel | None: + + """ + Gets one message if you know its message id. + :param mongo_conn: The instance of the database connector to use for the operation. + :param message_id: The id of the message that needs to be read. + :return: The contents of that one message in a structured format. + """ + + # We fetch the whole payload of that one message: + record = await mongo_conn.find_one( + collection = self.MESSAGES_COLLECTION, + filter = {"_id": ObjectId(message_id)}, + raise_exception = True + ) + + # If no such message was found: + if record is None: return None + + # If a record was found, + # we return it as our data model: + return CoreMessageModel(**record) + + # ┏┓┳┓┳┳┳┓ ┳┳ ┓ + # ┃ ┣┫┃┃┃┃ ━━ ┃┃┏┓┏┫┏┓╋┏┓ + # ┗┛┛┗┗┛┻┛ ┗┛┣┛┗┻┗┻┗┗ + # ┛ + + # We don't support updating messages themselves, + # but we will allow updating fields like tags, marking as read or unread, etc. + + async def update_tags( + self, + mongo_conn: AsyncMongo, + message_id: ObjectId | str, + unset_tags: List[str] = None, + set_tags: List[str] = None + ) -> bool: + + """ + Updates the tags on one message. The tags to remove are processed first, the ones to add are processed later. + :param mongo_conn: The instance of the database connector to use for the operation. + :param message_id: The id of the message that needs to be read. + :param unset_tags: The tags to remove from the message. + :param set_tags: The tags to add to the message. + :return: True if the update was successful, else False. + """ + + # Update the tags: + return await mongo_conn.update_one( + collection = self.MESSAGES_COLLECTION, + filter = {"_id": ObjectId(message_id)}, + update = [{ + "$set": { + "tags": { + "$let": { + "vars": { + "removed_tags": { + "$setDifference": [ + "$tags", + unset_tags + ] + } + }, + "in": { + "$setUnion": [ + "$$removed_tags", + set_tags + ] + } + } + } + } + }], + raise_exception = True + ) + + # ┏┓┳┓┳┳┳┓ ┳┓ ┓ + # ┃ ┣┫┃┃┃┃ ━━ ┃┃┏┓┃┏┓╋┏┓ + # ┗┛┛┗┗┛┻┛ ┻┛┗ ┗┗ ┗┗ + + # No support whatsoever for deleting messages. + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/controllers_v2/core/payment.py b/controllers_v2/core/payment.py new file mode 100644 index 0000000..c4dbb05 --- /dev/null +++ b/controllers_v2/core/payment.py @@ -0,0 +1,403 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Monday, 16th Dec., 2024 + + OBJECTIVE: + + To handle all payments from one place. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# For Quart: +from quart import current_app + +# 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, AsyncMongoStorage + +# Base model: +from controllers.base import BaseModel + +# Data models: +from models.core.auth_token import CoreAuthTokenModel +from models.core.payment import CorePaymentModel, PaymentEvent +from models.core.user import CoreUserInfoModel + +# To work with MongoDB: +from bson import ObjectId +from pymongo import InsertOne, UpdateOne, ReplaceOne + +# To work with datatypes: +from typing import Literal, List, Dict, Any + +# To make deep-copies: +import copy + +# To work with base-64 encoding: +import base64 + +# To work with date and time: +import datetime + +# For asynchronous activities: +import asyncio + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class CorePaymentController(BaseModel): + + # ┏┓┓ ┓┏ + # ┃ ┃┏┓┏┏ ┃┃┏┓┏┓┏ + # ┗┛┗┗┻┛┛ ┗┛┗┻┛ ┛ + + # For MongoDB: + PAYMENTS_COLLECTION = "_payments" + + # ┏┓┳┓┳┳┳┓ ┏┓ + # ┃ ┣┫┃┃┃┃ ━━ ┃ ┏┓┏┓┏┓╋┏┓ + # ┗┛┛┗┗┛┻┛ ┗┛┛ ┗ ┗┻┗┗ + + async def init( + self, + mongo_conn: AsyncMongo, + payment: CorePaymentModel + ) -> ObjectId: + + """ + Simply insert one payment document into the database. + :param mongo_conn: The instance of the database connector to use for the operation. + :param payment: The payment whose record needs to be saved in the database. + :return: The object id of the inserted document. + """ + + # Receive te JSON: + payment_json = payment.model_dump() + payment_json.pop("_id") + + # Simply insert the document: + return await mongo_conn.insert_one( + collection = self.PAYMENTS_COLLECTION, + document = payment_json, + raise_exception = True + ) + + # ┏┓┳┓┳┳┳┓ ┳┓ • + # ┃ ┣┫┃┃┃┃ ━━ ┣┫┏┓╋┏┓┓┏┓┓┏┏┓ + # ┗┛┛┗┗┛┻┛ ┛┗┗ ┗┛ ┗┗ ┗┛┗ + + async def count_payments( + self, + mongo_conn: AsyncMongo, + token_ids: List[ObjectId | str], + additional_filter: dict = None + ) -> int: + + """ + Just counts the no. of payment records that match a given set of conditions. + :param mongo_conn: The instance of the database connector to use for the operation. + :param token_ids: The token ids of the accounts from which these payment details must be fetched. + :param additional_filter: Any addition filters to use. + :return: The no. of payment records that match the given conditions. + """ + + # Prepare the filter: + if not isinstance(token_ids, list): token_ids = [token_ids] + token_ids = [ObjectId(t) for t in token_ids] + filter_json = {"tokenId": {"$in": token_ids}} + if additional_filter: + for k, v in additional_filter.items(): + filter_json[k] = v + + # Get the count of the documents that match the criteria: + count = await mongo_conn.count( + collection = self.PAYMENTS_COLLECTION, + filter = filter_json, + raise_exception = True + ) + + # Done here: + return count + + async def get_payment_previews( + self, + mongo_conn: AsyncMongo, + token_ids: List[ObjectId | str], + limit: int = 100, + skip: int = 0, + additional_filter: dict = None + ) -> List[CorePaymentModel] | None: + + """ + Fetches many payment details in one call, but just their previews. + :param mongo_conn: The instance of the database connector to use for the operation. + :param token_ids: The token ids of the accounts from which these messages must be fetched. + :param limit: The max. no. of payment details to retrieve in this call. + :param skip: The no. of initial payment details to skip. Useful for pagination. + :param additional_filter: Any addition filters to use. + :return: The list of payments (as the payments model). This list can be empty. + """ + + # Prepare the filter: + if not isinstance(token_ids, list): token_ids = [token_ids] + token_ids = [ObjectId(t) for t in token_ids] + filter_json = {"tokenId": {"$in": token_ids}} + if additional_filter: + for k, v in additional_filter.items(): + filter_json[k] = v + + # We fetch the messages that are identified by a specific token id, + # with the specified fetching limits, while enforcing the sorting condition: + records = await mongo_conn.find_many( + collection = self.PAYMENTS_COLLECTION, + filter = filter_json, + projection = {"events": False}, + limit = limit, + skip = skip, + sort = {"ts": -1}, + raise_exception = True + ) + + # Convert the fetched records to instances of the data model and return: + for record in records: record["events"] = [] + return [CorePaymentModel(**record) for record in records] + + async def get_payment( + self, + mongo_conn: AsyncMongo, + payment_id: ObjectId | str, + ) -> CorePaymentModel | None: + + """ + Gets one payment detail if you know its payment id. + :param mongo_conn: The instance of the database connector to use for the operation. + :param payment_id: The id of the payment detail that needs to be read. + :return: The contents of that one payment detail in a structured format. + """ + + # We fetch the whole payload of that one message: + record = await mongo_conn.find_one( + collection = self.PAYMENTS_COLLECTION, + filter = {"_id": ObjectId(payment_id)}, + raise_exception = True + ) + + # If no such message was found: + if record is None: return None + + # If a record was found, + # we return it as our data model: + return CorePaymentModel(**record) + + # ┏┓┳┓┳┳┳┓ ┳┳ ┓ + # ┃ ┣┫┃┃┃┃ ━━ ┃┃┏┓┏┫┏┓╋┏┓ + # ┗┛┛┗┗┛┻┛ ┗┛┣┛┗┻┗┻┗┗ + # ┛ + + # We don't support updating payments themselves, + # but we will allow updating fields like tags, adding events, etc. + + async def add_event_by_payment_id( + self, + mongo_conn: AsyncMongo, + payment_id: ObjectId | str, + event: PaymentEvent, + client_reference_id: str = None + ) -> bool: + + """ + Add an event to an existing record of a payment detail. + :param mongo_conn: The instance of the database connector to use for the operation. + :param payment_id: The id of the payment detail that needs to be read. + :param event: The event that occurred. This will typically be generated by the third-party client. + :param client_reference_id: The way the client identifies this payment. You need to pass this only on the first + event. Typically, when you initiate the payment request. + :return: True if successfully noted, else False. + """ + + # Prepare the update document: + update_json = { + "$push": { + "events": event.model_dump() + }, + "$set": { + "lastEventTs": event.eventTs, + "lastEventMessage": event.message, + "lastPaymentStatus": event.paymentStatus, + } + } + if client_reference_id: update_json["$set"]["clientPaymentReferenceId"] = client_reference_id + + # Try to update the existing record: + return await mongo_conn.update_one( + collection = self.PAYMENTS_COLLECTION, + filter = {"_id": ObjectId(payment_id)}, + update = update_json, + upsert = False, + raise_exception = True + ) + + async def add_event_by_client_reference_id( + self, + mongo_conn: AsyncMongo, + client_reference_id: str, + event: PaymentEvent, + ) -> bool: + + """ + Add an event to an existing record of a payment detail. + :param mongo_conn: The instance of the database connector to use for the operation. + :param event: The event that occurred. This will typically be generated by the third-party client. + :param client_reference_id: The way the client identifies this payment. You need to pass this only on the first + event. Typically, when you initiate the payment request. + :return: True if successfully noted, else False. + """ + + # Prepare the update document: + update_json = { + "$push": { + "events": event.model_dump() + }, + "$set": { + "lastEventTs": event.eventTs, + "lastEventMessage": event.message, + "lastPaymentStatus": event.paymentStatus + } + } + if client_reference_id: update_json["$set"]["clientPaymentReferenceId"] = client_reference_id + + # Try to update the existing record: + return await mongo_conn.update_one( + collection = self.PAYMENTS_COLLECTION, + filter = {"clientPaymentReferenceId": client_reference_id}, + update = update_json, + upsert = False, + raise_exception = True + ) + + async def update_tags( + self, + mongo_conn: AsyncMongo, + payment_id: ObjectId | str, + unset_tags: List[str] = None, + set_tags: List[str] = None + ) -> bool: + + """ + Updates the tags on one payment. The tags to remove are processed first, the ones to add are processed later. + :param mongo_conn: The instance of the database connector to use for the operation. + :param payment_id: The id of the payment detail that needs to be read. + :param unset_tags: The tags to remove from the payment record. + :param set_tags: The tags to add to the payment record. + :return: True if the update was successful, else False. + """ + + # Update the tags: + return await mongo_conn.update_one( + collection = self.PAYMENTS_COLLECTION, + filter = {"_id": ObjectId(payment_id)}, + update = [{ + "$set": { + "tags": { + "$let": { + "vars": { + "removed_tags": { + "$setDifference": [ + "$tags", + unset_tags + ] + } + }, + "in": { + "$setUnion": [ + "$$removed_tags", + set_tags + ] + } + } + } + } + }], + raise_exception = True + ) + + # ┏┓┳┓┳┳┳┓ ┳┓ ┓ + # ┃ ┣┫┃┃┃┃ ━━ ┃┃┏┓┃┏┓╋┏┓ + # ┗┛┛┗┗┛┻┛ ┻┛┗ ┗┗ ┗┗ + + # No support whatsoever for deleting messages. + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/controllers_v2/sms/__init__.py b/controllers_v2/sms/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/controllers_v2/sms/all_sms.py b/controllers_v2/sms/all_sms.py new file mode 100644 index 0000000..90a6690 --- /dev/null +++ b/controllers_v2/sms/all_sms.py @@ -0,0 +1,276 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 19th Dec., 2024 + + OBJECTIVE: + + To handle all SMS related behaviour for Nimbus It's service from one place. + This service is for India only. + + 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.date_time import date_time +from utils_v2.database.async_mongo_v2 import AsyncMongo +from utils_v2.cache.async_redis_cache_v2 import AsyncRedisCache + +# Controllers: +from controllers_v2.sms.base import SMSController + +# Models: +from models.core.auth_token import CoreAuthTokenModel +from models.core.message import CoreMessageModel +from models.api.sms.send import ( + NimbusSMSIndiaMessage, + SMSSendOneResult, + SMSSendManyResults +) + +# SMS Clients: +from utils_v2.sms.india.nimbus.controllers.async_nimbus import AsyncNimbusSMS + +# To work with datatypes: +from typing import List, Any + +# To make HTTP requests: +import httpx + +# For asynchronous activities: +import asyncio + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class NimbusSMSIndiaController(SMSController): + + # ┏┓ + # ┃ ┏┓┏┓┏╋┏┓┓┏┏╋┏┓┏┓ + # ┗┛┗┛┛┗┛┗┛ ┗┻┗┗┗┛┛ + + def __init__( + self, + cache: AsyncRedisCache = None, + http_client: httpx.AsyncClient = None, + alert_url: str = None, + debug: bool = True, + debug_prefix: str = "Nimbus SMS (C) | ", + debug_only_errors: bool = True + ): + + """ + This is the foundational controller for all SMS services. This is built on top of the core message controller, + and, in turn, all individual SMS client controllers must be built on top of this. + :param cache: The object to use for caching results from database calls. + :param http_client: The HTTP client + :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. + """ + + # Invoke the parent's constructor: + super().__init__( + cache = cache, + alert_url = alert_url, + http_client = http_client, + base_filter = {"client": "nimbusSmsIndia"}, + debug = debug, + debug_prefix = debug_prefix, + debug_only_errors = debug_only_errors + ) + + # ┏┓┳┳┓┏┓ ┏┓ ┓• + # ┗┓┃┃┃┗┓ ┗┓┏┓┏┓┏┫┓┏┓┏┓ + # ┗┛┛ ┗┗┛ ┗┛┗ ┛┗┗┻┗┛┗┗┫ + # ┛ + + async def send_one_sms( + self, + mongo_data_conn: AsyncMongo, + auth_token: CoreAuthTokenModel, + client: AsyncNimbusSMS, + message: NimbusSMSIndiaMessage, + tags: List[Any] + ) -> SMSSendOneResult: + + """ + Use this to send one SMS. There are just 2 steps here - send the SMS, and store its details in the database. + :param mongo_data_conn: The database connection to use to perform this task. + :param auth_token: The auth token that will be used to send this message. + :param client: The third-party SMS client to use to send this message. + :param message: The actual message that needs to be sent. + :param tags: Any tags to attach with this SMS for filtering when querying in the listing service. + :return: The structured result of sending one message. + """ + + # Send the SMS: + client_response = await client.send_sms( + recipient_number = message.recipientNo, + message = message.text, + template_id = message.templateId + ) + + # Convert the format of the SMS client's response to the core message model. + sent_message_model = CoreMessageModel( + ts = client_response.ts, + syncTs = date_time.get_current_utc_date_time(as_string = False), + tokenId = auth_token.authTokenId, + serviceType = auth_token.serviceType, + client = auth_token.client, + clientMessageId = client_response.messageId, + clientThreadId = message.recipientNo, + isSent = True, + isBroadcast = False, + sentSuccessfully = client_response.success, + sender = None, + recipient = message.recipientNo, + chat = None, + message = client_response.model_dump(), + snippet = message.text, + aiSnippet = None, + tags = list(set(tags + ["SMS", "Nimbus SMS", "India"])) + ) + + # Save the result to the database: + message_id = await self.save_one_message( + mongo_data_conn = mongo_data_conn, + message = sent_message_model + ) + self._printer(message_id, client_response.success) + + # Done here: + success = True if client_response.success and message_id else False + return SMSSendOneResult( + success = success, + message = "SMS sent successfully." if client_response.success else "SMS sending failed.", + smsMessage = message + ) + + async def send_many_sms( + self, + mongo_data_conn: AsyncMongo, + auth_token: CoreAuthTokenModel, + messages: List[NimbusSMSIndiaMessage], + tags: List[Any] + ) -> SMSSendManyResults: + + """ + Use this to send multiple SMS messages. This method just calls the individual SMS sending method for every + individual message, and then aggregates the results. + :param mongo_data_conn: The database connection to use to perform this task. + :param auth_token: The auth token that will be used to send this message. + :param messages: The list of messages to send out. + :param tags: Any tags to attach with these SMS for filtering when querying in the listing service. The same tags + will be applied to all messages. Do not call this method if you need to have different tags for all of them. + :return: The structured result of sending many SMS messages. + """ + + # Start with a blank variable: + cumulative_results = SMSSendManyResults() + + # Make the client from the auth-token: + client = AsyncNimbusSMS( + entity_id = auth_token.auth["entityId"], + sender_id = auth_token.auth["senderId"], + user_id = auth_token.auth["userId"], + api_key = auth_token.auth["apiKey"], + http_client = self._http_client + ) + + # Create and fire all the SMS-sending tasks: + tasks = [ + self.send_one_sms( + mongo_data_conn = mongo_data_conn, + auth_token = auth_token, + client = client, + message = message, + tags = tags + ) + for message in messages + ] + individual_results = await asyncio.gather(*tasks) + + # Prepare the final result: + for result in individual_results: + if result.success: cumulative_results.successCount += 1 + else: cumulative_results.failureCount += 1 + cumulative_results.totalCount += 1 + cumulative_results.smsMessages.append(result.smsMessage) + cumulative_results.message = f"{cumulative_results.successCount}/{cumulative_results.totalCount} SMS sent." + + # Done here: + return cumulative_results + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/controllers_v2/sms/base.py b/controllers_v2/sms/base.py new file mode 100644 index 0000000..1b4fae9 --- /dev/null +++ b/controllers_v2/sms/base.py @@ -0,0 +1,434 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Friday, 13th Dec., 2024 + + OBJECTIVE: + + To handle all SMS related behaviour from one place. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# For Quart: +from quart import current_app + +# 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, AsyncMongoStorage + +# Base model: +from controllers.base import BaseModel + +# Data models: +from models.core.user import CoreUserInfoModel +from models.core.auth_token import CoreAuthTokenModel +from models.core.message import CoreMessageModel +from models.api.sms.send import ( + SMSSendRequestData, + NimbusSMSIndiaMessage, + SavvyBulkSMSKenyaMessage, + SMSSendManyResults +) + +# SMS Clients: +from utils_v2.sms.india.nimbus.controllers.async_nimbus import AsyncNimbusSMS +from utils_v2.sms.kenya.savvy_bulk_sms.controllers.async_savvy_bulk_sms import AsyncSavvyBulkSMS +from utils_v2.sms.models.sms_message import SentSMSMessageModel + +# To work with MongoDB: +from bson import ObjectId +from pymongo import InsertOne + +# To work with datatypes: +from typing import Literal, List, Dict, Any + +# To make API calls: +import httpx + +# For asynchronous activities: +import asyncio + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class SMSController: + + # ┏┓┓ ┓┏ + # ┃ ┃┏┓┏┏ ┃┃┏┓┏┓┏ + # ┗┛┗┗┻┛┛ ┗┛┗┻┛ ┛ + + pass + + # ┓┏ ┓ + # ┣┫┏┓┃┏┓┏┓┏┓┏ + # ┛┗┗ ┗┣┛┗ ┛ ┛ + # ┛ + + pass + + # ┏┓ ┓ + # ┣┫┓┏╋┣┓ + # ┛┗┗┻┗┛┗ + + @staticmethod + async def set_token_direct( + db_conn: AsyncMySQL, + mongo_conn: AsyncMongo, + auth_token: CoreAuthTokenModel, + session_token: str = None + ) -> bool: + + # Start by assuming failure: + success = False + + # Get a token id: + token_key = await current_app.core_auth_token_controller.get_token_key( + db_conn = db_conn, + mongo_conn = mongo_conn, + auth_token = auth_token, + token_notes = {}, + session_token = session_token + ) + + # Immediately save the details against that token id: + success = await current_app.core_auth_token_controller.set_token( + db_conn = db_conn, + mongo_conn = mongo_conn, + token_key = token_key, + auth_token = auth_token, + token_notes = {}, + session_token = session_token + ) + + # Done here: + return success + + @staticmethod + async def get_token( + mongo_conn: AsyncMongo, + token_key: ObjectId | str = None, + ) -> CoreAuthTokenModel | None: + + # Simply call the core model: + return await current_app.core_auth_token_controller.get_token_from_key( + mongo_conn = mongo_conn, + token_key = token_key + ) + + # ┏┓ ┓ + # ┗┓┏┓┏┓┏┫ + # ┗┛┗ ┛┗┗┻ + + @staticmethod + async def __send_from_nimbus_sms_india( + http_client: httpx.AsyncClient, + auth_token: CoreAuthTokenModel, + messages: List[NimbusSMSIndiaMessage], + tags: List[Any] + ) -> SMSSendManyResults: + + # Start with a blank variable: + send_results = SMSSendManyResults() + + # Initialize the third-party client: + client = AsyncNimbusSMS( + entity_id = auth_token.auth["entityId"], + sender_id = auth_token.auth["senderId"], + user_id = auth_token.auth["userId"], + api_key = auth_token.auth["apiKey"], + http_client = http_client + ) + + # Iterate over all the messages you need to send: + for message in messages: + + # Send the SMS and return the response: + client_response = await client.send_sms( + recipient_number = message.recipientNo, + message = message.text, + template_id = message.templateId + ) + + # Note down the results: + send_results.totalCount += 1 + if client_response.success: send_results.successCount += 1 + else: send_results.failureCount += 1 + send_results.smsMessages.append(CoreMessageModel( + ts = client_response.ts, + syncTs = date_time.get_current_utc_date_time(as_string = False), + tokenId = auth_token.authTokenId, + serviceType = auth_token.serviceType, + client = auth_token.client, + clientMessageId = client_response.messageId, + clientThreadId = message.recipientNo, + isSent = True, + isBroadcast = False, + sentSuccessfully = client_response.success, + sender = None, + recipient = message.recipientNo, + chat = None, + message = client_response.model_dump(), + snippet = message.text, + aiSnippet = None, + tags = list(set(tags + ["SMS", "Nimbus SMS", "India"])) + )) + + # Done here: + return send_results + + @staticmethod + async def __send_from_savvy_bulk_sms_kenya( + http_client: httpx.AsyncClient, + auth_token: CoreAuthTokenModel, + messages: List[SavvyBulkSMSKenyaMessage], + tags: List[Any] + ) -> SMSSendManyResults: + + # Start with a blank variable: + send_results = SMSSendManyResults() + + # Initialize the third-party client: + client = AsyncSavvyBulkSMS( + partner_id = auth_token.auth["partnerId"], + short_code = auth_token.auth["shortCode"], + api_key = auth_token.auth["apiKey"], + http_client = http_client + ) + + # Iterate over all the messages you need to send: + for message in messages: + + # Send the SMS and return the response: + client_response = await client.send_sms( + recipient_number = message.recipientNo, + message = message.text + ) + + # Note down the results: + send_results.totalCount += 1 + if client_response.success: send_results.successCount += 1 + else: send_results.failureCount += 1 + send_results.smsMessages.append(CoreMessageModel( + ts = client_response.ts, + syncTs = date_time.get_current_utc_date_time(as_string = False), + tokenId = auth_token.authTokenId, + serviceType = auth_token.serviceType, + client = auth_token.client, + clientMessageId = client_response.messageId, + clientThreadId = message.recipientNo, + isSent = True, + isBroadcast = False, + sentSuccessfully = client_response.success, + sender = None, + recipient = message.recipientNo, + chat = None, + message = client_response.model_dump(), + snippet = message.text, + aiSnippet = None, + tags = list(set(tags + ["SMS", "Savvy Bulk SMS", "Kenya"])) + )) + + # Done here: + return send_results + + async def send( + self, + mongo_conn: AsyncMongo, + http_client: httpx.AsyncClient, + auth_token: CoreAuthTokenModel, + messages: List[NimbusSMSIndiaMessage | SavvyBulkSMSKenyaMessage], + tags: List[Any] + ) -> SMSSendManyResults: + + # Start by assuming failure: + send_results = SMSSendManyResults() + + # Now we route the message to the appropriate client: + match auth_token.client: + case "nimbusSmsIndia": + send_results = await self.__send_from_nimbus_sms_india( + http_client = http_client, + auth_token = auth_token, + messages = messages, + tags = tags + ) + case "savvyBulkSmsKenya": + send_results = await self.__send_from_savvy_bulk_sms_kenya( + http_client = http_client, + auth_token = auth_token, + messages = messages, + tags = tags + ) + case _: + send_results.message = f"invalid client {auth_token.client}" + + # Save the results to MongoDB: + tasks = [] + for sms in send_results.smsMessages: + message_json = sms.model_dump() + message_json.pop("_id", None) + tasks.append(current_app.core_message_controller.insert( + mongo_conn = mongo_conn, + message = message_json + )) + results = await asyncio.gather(*tasks) + + # Done here: + send_results.message = f"{send_results.successCount}/{send_results.totalCount} message(s) sent" + return send_results + + # ┓ • ┏┓ ┏┓ ┳┳┓ + # ┃ ┓┏╋ ┣╋ ┃┓┏┓╋ ┃┃┃┏┓┏┏┏┓┏┓┏┓┏ + # ┗┛┗┛┗ ┗┻ ┗┛┗ ┗ ┛ ┗┗ ┛┛┗┻┗┫┗ ┛ + # ┛ + + # These are simply for retrieving sms messages. + # You need to already have them saved to the database. + + # @staticmethod + # async def list_messages( + # mongo_conn: AsyncMongo, + # token_ids: List[ObjectId | str], + # limit: int = 100, + # skip: int = 0, + # additional_filter: dict = None + # ) -> List[CoreMessageModel] | None: + # + # # regardless of what additional filter is provided from outside, + # # we add a mail-selecting filter here: + # if additional_filter is None: additional_filter = {} + # additional_filter["serviceType"] = "sms" + # + # # Simply call the core model: + # return await current_app.core_message_controller.get_message( + # mongo_conn = mongo_conn, + # token_ids = token_ids, + # limit = limit, + # skip = skip, + # additional_filter = additional_filter + # ) + # + # @staticmethod + # async def get_one_mail( + # mongo_conn: AsyncMongo, + # token_id: ObjectId | str, + # message_id: ObjectId | str + # ) -> CoreMessageModel | None: + # + # # Simply call the core model: + # return await current_app.core_message_controller.get_message( + # mongo_conn = mongo_conn, + # token_id = token_id, + # message_id = message_id + # ) + + # ┳┳ ┓ + # ┃┃┏┓┏┫┏┓╋┏┓ + # ┗┛┣┛┗┻┗┻┗┗ + # ┛ + + @staticmethod + async def update_tags( + mongo_conn: AsyncMongo, + token_id: ObjectId | str, + message_id: ObjectId | str, + unset_tags: List[str] = None, + set_tags: List[str] = None + ) -> bool: + + # Simply call the core model: + return await current_app.core_message_controller.update_tags( + mongo_conn = mongo_conn, + token_id = token_id, + message_id = message_id, + unset_tags = unset_tags, + set_tags = set_tags + ) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass + + # from utils_v2.string import json + # + # file_options = [ + # r"/home/developer/Downloads/recursive parts parse - 20241210.json", + # r"/home/developer/Downloads/recursive parts parse (no attachment) - 20241210.json", + # ] + # + # raw_mail_json = json.from_file(file_options[1]) + # print("FROM FILE:", json.to_string(raw_mail_json["payload"])) + # print("\n\n---------\n\n") + # mail_controller = MailController() + # print(json.to_string(mail_controller.drop_attachments(raw_mail_json["payload"]))) diff --git a/controllers_v2/sms/nimbus_sms_india.py b/controllers_v2/sms/nimbus_sms_india.py new file mode 100644 index 0000000..8965ae0 --- /dev/null +++ b/controllers_v2/sms/nimbus_sms_india.py @@ -0,0 +1,197 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Friday, 13th Dec., 2024 + + OBJECTIVE: + + To handle all SMS related behaviour from one place. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# For Quart: +from quart import current_app + +# 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, AsyncMongoStorage +from utils_v2.cache.async_redis_cache_v2 import AsyncRedisCache + +# Controllers: +from controllers_v2.core.message import CoreMessageController + +# Models: +from models.core.auth_token import CoreAuthTokenModel +from models.core.message import CoreMessageModel +from models.core.user import CoreUserInfoModel +from models.api.sms.send import ( + SMSSendRequestData, + NimbusSMSIndiaMessage, + SavvyBulkSMSKenyaMessage, + SMSSendManyResults +) + +# To work with MongoDB: +from bson import ObjectId +from pymongo import InsertOne, UpdateOne, ReplaceOne + +# To work with datatypes: +from typing import Literal, List, Dict, Any + +# To make HTTP requests: +import httpx + +# To make deep-copies: +import copy + +# To work with base-64 encoding: +import base64 + +# To work with date and time: +import datetime + +# For asynchronous activities: +import asyncio + +# To make abstract classes: +from abc import ABC, abstractmethod + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class SMSController(CoreMessageController, ABC): + + # ┏┓ + # ┃ ┏┓┏┓┏╋┏┓┓┏┏╋┏┓┏┓ + # ┗┛┗┛┛┗┛┗┛ ┗┻┗┗┗┛┛ + + def __init__( + self, + cache: AsyncRedisCache = None, + http_client: httpx.AsyncClient = None, + alert_url: str = None, + base_filter: dict = None, + debug: bool = True, + debug_prefix: str = "SMS (C) | ", + debug_only_errors: bool = True + ): + + """ + This is the foundational controller for all SMS services. This is built on top of the core message controller, + and, in turn, all individual SMS client controllers must be built on top of this. + :param cache: The object to use for caching results from database calls. + :param http_client: The HTTP client + :param base_filter: The basic filter that will be applied to all fetching/updating queries. WARNING: THE BASE + FILTER WILL ALWAYS BE APPLIED AUTOMATICALLY. SET THIS UP WISELY. + :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. + """ + + # Prepare the combined base filter: + sms_filter = {"serviceType": "sms"} + if base_filter: + for k, v in base_filter.items(): sms_filter[k] = v + + # Invoke the parent's constructor: + super().__init__( + cache = cache, + alert_url = alert_url, + http_client = http_client, + base_filter = sms_filter, + debug = debug, + debug_prefix = debug_prefix, + debug_only_errors = debug_only_errors + ) + + # ┏┓┳┳┓┏┓ ┏┓ ┓• + # ┗┓┃┃┃┗┓ ┗┓┏┓┏┓┏┫┓┏┓┏┓ + # ┗┛┛ ┗┗┛ ┗┛┗ ┛┗┗┻┗┛┗┗┫ + # ┛ + + @abstractmethod + async def send( + self, + mongo_conn: AsyncMongo, + auth_token: CoreAuthTokenModel, + messages: List[NimbusSMSIndiaMessage | SavvyBulkSMSKenyaMessage], + tags: List[Any] + ) -> SMSSendManyResults: + + pass + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/controllers_v2/sms/savvy_bulk_sms_kenya.py b/controllers_v2/sms/savvy_bulk_sms_kenya.py new file mode 100644 index 0000000..ccf139b --- /dev/null +++ b/controllers_v2/sms/savvy_bulk_sms_kenya.py @@ -0,0 +1,273 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 19th Dec., 2024 + + OBJECTIVE: + + To handle all SMS related behaviour for Nimbus It's service from one place. + + 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.date_time import date_time +from utils_v2.database.async_mongo_v2 import AsyncMongo +from utils_v2.cache.async_redis_cache_v2 import AsyncRedisCache + +# Controllers: +from controllers_v2.sms.base import SMSController + +# Models: +from models.core.auth_token import CoreAuthTokenModel +from models.core.message import CoreMessageModel +from models.api.sms.send import ( + NimbusSMSIndiaMessage, + SMSSendOneResult, + SMSSendManyResults +) + +# SMS Clients: +from utils_v2.sms.india.nimbus.controllers.async_nimbus import AsyncNimbusSMS + +# To work with datatypes: +from typing import List, Any + +# To make HTTP requests: +import httpx + +# For asynchronous activities: +import asyncio + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** CLASSES *** +# ***** **** +# ***************************************************************************************************************** + + +class NimbusSMSIndiaController(SMSController): + + # ┏┓ + # ┃ ┏┓┏┓┏╋┏┓┓┏┏╋┏┓┏┓ + # ┗┛┗┛┛┗┛┗┛ ┗┻┗┗┗┛┛ + + def __init__( + self, + cache: AsyncRedisCache = None, + http_client: httpx.AsyncClient = None, + alert_url: str = None, + debug: bool = True, + debug_prefix: str = "SMS (C) | ", + debug_only_errors: bool = True + ): + + """ + This is the foundational controller for all SMS services. This is built on top of the core message controller, + and, in turn, all individual SMS client controllers must be built on top of this. + :param cache: The object to use for caching results from database calls. + :param http_client: The HTTP client + :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. + """ + + # Invoke the parent's constructor: + super().__init__( + cache = cache, + alert_url = alert_url, + http_client = http_client, + base_filter = {"client": "nimbusSmsIndia"}, + debug = debug, + debug_prefix = debug_prefix, + debug_only_errors = debug_only_errors + ) + + # ┏┓┳┳┓┏┓ ┏┓ ┓• + # ┗┓┃┃┃┗┓ ┗┓┏┓┏┓┏┫┓┏┓┏┓ + # ┗┛┛ ┗┗┛ ┗┛┗ ┛┗┗┻┗┛┗┗┫ + # ┛ + + async def send_one_sms( + self, + mongo_data_conn: AsyncMongo, + auth_token: CoreAuthTokenModel, + client: AsyncNimbusSMS, + message: NimbusSMSIndiaMessage, + tags: List[Any] + ) -> SMSSendOneResult: + + """ + Use this to send one SMS. There are just 2 steps here - send the SMS, and store its details in the database. + :param mongo_data_conn: The database connection to use to perform this task. + :param auth_token: The auth token that will be used to send this message. + :param client: The third-party SMS client to use to send this message. + :param message: The actual message that needs to be sent. + :param tags: Any tags to attach with this SMS for filtering when querying in the listing service. + :return: The structured result of sending one message. + """ + + # Send the SMS: + client_response = await client.send_sms( + recipient_number = message.recipientNo, + message = message.text, + template_id = message.templateId + ) + + # Convert the format of the SMS client's response to the core message model. + sent_message_model = CoreMessageModel( + ts = client_response.ts, + syncTs = date_time.get_current_utc_date_time(as_string = False), + tokenId = auth_token.authTokenId, + serviceType = auth_token.serviceType, + client = auth_token.client, + clientMessageId = client_response.messageId, + clientThreadId = message.recipientNo, + isSent = True, + isBroadcast = False, + sentSuccessfully = client_response.success, + sender = None, + recipient = message.recipientNo, + chat = None, + message = client_response.model_dump(), + snippet = message.text, + aiSnippet = None, + tags = list(set(tags + ["SMS", "Nimbus SMS", "India"])) + ) + + # Save the result to the database: + message_id = await self.save_one_message( + mongo_data_conn = mongo_data_conn, + message = sent_message_model + ) + + # Done here: + success = True if client_response.success and message_id else False + return SMSSendOneResult( + success = success, + message = "SMS sent successfully." if client_response.success else "SMS sending failed.", + smsMessage = sent_message_model + ) + + async def send_many_sms( + self, + mongo_data_conn: AsyncMongo, + auth_token: CoreAuthTokenModel, + messages: List[NimbusSMSIndiaMessage], + tags: List[Any] + ) -> SMSSendManyResults: + + """ + Use this to send multiple SMS messages. This method just calls the individual SMS sending method for every + individual message, and then aggregates the results. + :param mongo_data_conn: The database connection to use to perform this task. + :param auth_token: The auth token that will be used to send this message. + :param messages: The list of messages to send out. + :param tags: Any tags to attach with these SMS for filtering when querying in the listing service. The same tags + will be applied to all messages. Do not call this method if you need to have different tags for all of them. + :return: The structured result of sending many SMS messages. + """ + + # Start with a blank variable: + cumulative_results = SMSSendManyResults() + + # Make the client from the auth-token: + client = AsyncNimbusSMS( + entity_id = auth_token.auth["entityId"], + sender_id = auth_token.auth["senderId"], + user_id = auth_token.auth["userId"], + api_key = auth_token.auth["apiKey"], + http_client = self._http_client + ) + + # Create and fire all the SMS-sending tasks: + tasks = [ + self.send_one_sms( + mongo_data_conn = mongo_data_conn, + auth_token = auth_token, + client = client, + message = message, + tags = tags + ) + for message in messages + ] + individual_results = await asyncio.gather(*tasks) + + # Prepare the final result: + for result in individual_results: + if result.success: cumulative_results.successCount += 1 + else: cumulative_results.failureCount += 1 + cumulative_results.totalCount += 1 + cumulative_results.smsMessages.append(result.smsMessage) + + # Done here: + return cumulative_results + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/models/api/sms/list.py b/models/api/sms/list.py new file mode 100644 index 0000000..1191c21 --- /dev/null +++ b/models/api/sms/list.py @@ -0,0 +1,157 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Tuesday, 3rd Dec., 2024. + + OBJECTIVE: + + To provide a structure to query the full payload of an email. + + 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 +from typing import Optional, Literal, List, Any + +# My utils: +from utils_v2.string import regex +from utils_v2.date_time import date_time + +# To work with date and time: +import datetime + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# RegEx Patterns: +REGEX_SESSION_TOKEN = r"^[a-f0-9]{8}-[a-f0-9]{4}-[1-5][a-f0-9]{3}-[89ab][a-f0-9]{3}-[a-f0-9]{12}$" + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +class MailListRequestHeaders(BaseModel): + + sessionToken: str = Field( + description = "the session token of the user who is requesting the service", + pattern = REGEX_SESSION_TOKEN, + frozen = True, + alias = "X-Session-Token" + ) + + # ┏┓ ┏• + # ┃ ┏┓┏┓╋┓┏┓ + # ┗┛┗┛┛┗┛┗┗┫ + # ┛ + + class Config: + extra = "allow" + + def model_dump(self, *args, **kwargs): + return super().model_dump(*args, by_alias = True, **kwargs) + + +# --------------------------------------------------------------------------------------------------------------------- + + +class MailListRequestData(BaseModel): + + tokenKeys: str | List[str] = Field( + description = "the token identifier(s) that tell you which auth-tokens were used for fetching those messages", + frozen = True, + ) + + count: int = Field( + description = "the no. of mails to list", + default = 25, + ge = 1, + le = 500, + frozen = True + ) + + fromCount: int = Field( + description = "the no. of mails to skip before picking mails to list; useful for pagination", + ge = 0, + default = 0, + frozen = True + ) + + tags: List[Any] | None = Field( + description = "any no. of tags that you want to filter by", + default = None + ) + + # ┏┓ ┏• + # ┃ ┏┓┏┓╋┓┏┓ + # ┗┛┗┛┛┗┛┗┗┫ + # ┛ + + class Config: + extra = "forbid" + + # ┓┏ ┓• ┓ • + # ┃┃┏┓┃┓┏┫┏┓╋┓┏┓┏┓ + # ┗┛┗┻┗┗┗┻┗┻┗┗┗┛┛┗ + + @field_validator("tokenKeys", "tags", mode = "before") + def ensure_unique_list(cls, value): + if not isinstance(value, list): value = [value] + value = list(set(value)) + return value + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/models/api/sms/tags.py b/models/api/sms/tags.py new file mode 100644 index 0000000..dc2a39f --- /dev/null +++ b/models/api/sms/tags.py @@ -0,0 +1,139 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Friday, 13th Dec., 2024. + + OBJECTIVE: + + To provide a structure to work with the tags on mail messages. + + 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 +from typing import Optional, Literal, List, Any + +# My utils: +from utils_v2.string import regex +from utils_v2.date_time import date_time + +# To work with date and time: +import datetime + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# RegEx Patterns: +REGEX_SESSION_TOKEN = r"^[a-f0-9]{8}-[a-f0-9]{4}-[1-5][a-f0-9]{3}-[89ab][a-f0-9]{3}-[a-f0-9]{12}$" + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +class MailUpdateTagsRequestHeaders(BaseModel): + + sessionToken: str = Field( + description = "the session token of the user who is requesting the service", + pattern = REGEX_SESSION_TOKEN, + frozen = True, + alias = "X-Session-Token" + ) + + # ┏┓ ┏• + # ┃ ┏┓┏┓╋┓┏┓ + # ┗┛┗┛┛┗┛┗┗┫ + # ┛ + + class Config: + extra = "allow" + + def model_dump(self, *args, **kwargs): + return super().model_dump(*args, by_alias = True, **kwargs) + + +# --------------------------------------------------------------------------------------------------------------------- + + +class MailUpdateTagsRequestData(BaseModel): + + messageId: str = Field( + description = "the mail identifier (Mongo ObjectId) of the document that holds the mail", + frozen = True + ) + + unsetTags: List[Any] | None = Field( + description = "the list of tags to remove from the mail", + frozen = True, + default = None + ) + + setTags: List[Any] | None = Field( + description = "the list of tags to add to the mail", + frozen = True, + default = None + ) + + # ┏┓ ┏• + # ┃ ┏┓┏┓╋┓┏┓ + # ┗┛┗┛┛┗┛┗┗┫ + # ┛ + + class Config: + extra = "forbid" + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/utils_v2/api/async_quart.py b/utils_v2/api/async_quart.py index 1402e0d..a0113ef 100644 --- a/utils_v2/api/async_quart.py +++ b/utils_v2/api/async_quart.py @@ -26,7 +26,7 @@ DECORATORS. """ -import copy + # ***************************************************************************************************************** # ***** **** diff --git a/utils_v2/goog/controllers/gmail/__init__.py b/utils_v2/goog/controllers/gmail/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/utils_v2/goog/gmail/gmail_client.py b/utils_v2/goog/controllers/gmail/gmail_client.py similarity index 99% rename from utils_v2/goog/gmail/gmail_client.py rename to utils_v2/goog/controllers/gmail/gmail_client.py index 63e61d7..be458ed 100644 --- a/utils_v2/goog/gmail/gmail_client.py +++ b/utils_v2/goog/controllers/gmail/gmail_client.py @@ -51,7 +51,7 @@ from utils_v2.mail import mail_parser from utils_v2.goog.controllers.base import AsyncGoogleBase from utils_v2.goog.models.auth_tokens import GoogleAuthTokens from utils_v2.goog.models.api_call import GoogleApiResponse -from utils_v2.goog.gmail.gmail_message import GMailMessage +from utils_v2.goog.controllers.gmail.gmail_message import GMailMessage # Related to Google: from google.auth.transport.requests import Request diff --git a/utils_v2/goog/gmail/gmail_message.py b/utils_v2/goog/controllers/gmail/gmail_message.py similarity index 95% rename from utils_v2/goog/gmail/gmail_message.py rename to utils_v2/goog/controllers/gmail/gmail_message.py index 62556bd..d6ea111 100644 --- a/utils_v2/goog/gmail/gmail_message.py +++ b/utils_v2/goog/controllers/gmail/gmail_message.py @@ -101,10 +101,10 @@ class GMailMessage: def __init__( self, from_email: str, - to_email: str, + to_email: str | List[str], subject: str, - cc_emails: List[str] = None, - bcc_emails: List[str] = None + cc_emails: str | List[str] = None, + bcc_emails: str | List[str] = None ): """ @@ -119,10 +119,10 @@ class GMailMessage: # Create the instance of the message: self.message = MIMEMultipart() self.message["From"] = from_email - self.message["To"] = to_email + self.message["To"] = ",".join(to_email if isinstance(to_email, list) else [to_email]) self.message["Subject"] = subject - if cc_emails: self.message["CC"] = ",".join(cc_emails) - if bcc_emails: self.message["BCC"] = ",".join(bcc_emails) + if cc_emails: self.message["CC"] = ",".join(cc_emails if isinstance(cc_emails, list) else [cc_emails]) + if bcc_emails: self.message["BCC"] = ",".join(bcc_emails if isinstance(bcc_emails, list) else [bcc_emails]) # Note down the values for accessing later: self.__from = from_email @@ -331,6 +331,6 @@ if __name__ == "__main__": """ ) my_mail.add_text("This is how you should pet her 👇") - my_mail.add_inline_image(r"../../../data/images/cat_petting.png") - my_mail.add_attachment(r"../../../data/pdf/sample_label.pdf") - print(my_mail.get_raw_message(as_base64 = True)) + my_mail.add_inline_image(r"../../../../data/images/cat_petting.png") + my_mail.add_attachment(r"../../../../data/pdf/sample_label.pdf") + print(my_mail.get_raw_message(as_base64 = False)) diff --git a/utils_v2/goog/gmail/json/messages/chatgpt_suggested_message_format.json b/utils_v2/goog/controllers/gmail/json/messages/chatgpt_suggested_message_format.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/chatgpt_suggested_message_format.json rename to utils_v2/goog/controllers/gmail/json/messages/chatgpt_suggested_message_format.json diff --git a/utils_v2/goog/gmail/json/messages/formatted_message_list_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/formatted_message_list_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/formatted_message_list_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/formatted_message_list_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_01_193669615aa33694_no_attachment_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_01_193669615aa33694_no_attachment_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_01_193669615aa33694_no_attachment_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_01_193669615aa33694_no_attachment_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_02_19367930033154ca_with_attachment_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_02_19367930033154ca_with_attachment_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_02_19367930033154ca_with_attachment_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_02_19367930033154ca_with_attachment_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_03_19365d0239aa1aaf_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_03_19365d0239aa1aaf_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_03_19365d0239aa1aaf_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_03_19365d0239aa1aaf_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_04_193640e74d93dc9b_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_04_193640e74d93dc9b_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_04_193640e74d93dc9b_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_04_193640e74d93dc9b_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_05_19363d6a7f50225a_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_05_19363d6a7f50225a_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_05_19363d6a7f50225a_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_05_19363d6a7f50225a_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_06_19363ba8e2582133_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_06_19363ba8e2582133_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_06_19363ba8e2582133_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_06_19363ba8e2582133_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_07_19362f9983e58ea6_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_07_19362f9983e58ea6_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_07_19362f9983e58ea6_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_07_19362f9983e58ea6_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_08_19360aabcb976cfc_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_08_19360aabcb976cfc_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_08_19360aabcb976cfc_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_08_19360aabcb976cfc_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_09_193604ad97c9fe62_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_09_193604ad97c9fe62_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_09_193604ad97c9fe62_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_09_193604ad97c9fe62_20241126.json diff --git a/utils_v2/goog/gmail/json/messages/raw_message_10_1935e92672e67f22_20241126.json b/utils_v2/goog/controllers/gmail/json/messages/raw_message_10_1935e92672e67f22_20241126.json similarity index 100% rename from utils_v2/goog/gmail/json/messages/raw_message_10_1935e92672e67f22_20241126.json rename to utils_v2/goog/controllers/gmail/json/messages/raw_message_10_1935e92672e67f22_20241126.json