From bef8f6ef0d702958f54977afd013461ba96d61e1 Mon Sep 17 00:00:00 2001 From: khushal Date: Tue, 3 Dec 2024 16:01:08 +0530 Subject: [PATCH] (20241203) Mails sync'ing, listing, and fetching done. --- api/blueprints/mail/list.py | 31 +++---- api/blueprints/mail/oauth_callback.py | 3 - api/blueprints/mail/retrieve.py | 106 +++++------------------- api/blueprints/mail/sync.py | 2 +- api/main.py | 24 +++++- models/behaviour/mail/retrieve.py | 59 ++++++++++++-- models/data/mail/list.py | 22 ++++- models/data/mail/retrieve.py | 113 ++------------------------ 8 files changed, 137 insertions(+), 223 deletions(-) diff --git a/api/blueprints/mail/list.py b/api/blueprints/mail/list.py index 8889864..3827b9e 100644 --- a/api/blueprints/mail/list.py +++ b/api/blueprints/mail/list.py @@ -90,7 +90,7 @@ import datetime # Related to Quart: -mail_retrieve_bp = Blueprint("mail_retrieve", __name__) +mail_list_bp = Blueprint("mail_list", __name__) # ***************************************************************************************************************** @@ -110,7 +110,7 @@ mail_retrieve_bp = Blueprint("mail_retrieve", __name__) # ***************************************************************************************************************** -@mail_retrieve_bp.record_once +@mail_list_bp.record_once def init(blueprint_setup_state): # This gets called when the blueprint is registered. @@ -121,7 +121,7 @@ def init(blueprint_setup_state): # --------------------------------------------------------------------------------------------------------------------- -@mail_retrieve_bp.route("/list/account/id", methods = ["GET"]) +@mail_list_bp.route("/list/account/id", 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") @@ -137,19 +137,19 @@ def init(blueprint_setup_state): @log_chain_to_mongo(attr_name = "logs_mongo") @should_not_be_under_maintenance(attr_name = "is_under_maintenance") @validate_input( - header_validator = lambda x: MailGetRequestHeaders(**x).model_dump(), - data_validator = lambda x: MailGetRequestData(**x) + header_validator = lambda x: MailListRequestHeaders(**x).model_dump(), + data_validator = lambda x: MailListByAccountIdRequestData(**x) ) @handle_cancelled_request() -async def get_one_mail( - inbound_headers: dict | MailGetRequestHeaders = None, - inbound_data: dict | MailGetRequestData = None, +async def list_mails_for_account_id( + inbound_headers: dict | MailListRequestHeaders = None, + inbound_data: dict | MailListByAccountIdRequestData = None, inbound_files: dict = None, **kwargs ): """ - Use this endpoint when the user wants to fetch one mail. + Use this endpoint when the user wants to fetch the list of mails. :param inbound_headers: auto-extracted by the decorators. :param inbound_data: auto-extracted by the decorators. :param inbound_files: auto-extracted by the decorators. @@ -165,16 +165,19 @@ async def get_one_mail( ) # Get the mail: - mail_data = await current_app.mail_retrieve_model.get_mail( + mails_list = await current_app.mail_retrieve_model.list_for_account_identifier( mongo_conn = current_app.data_mongo, - mail_identifier = inbound_data.mailId + account_identifier = inbound_data.accountId, + limit = inbound_data.count, + skip = inbound_data.fromCount ) # Done here: return ResponseModel( - status_code = StatusCodes.OK if mail_data else StatusCodes.FAILED, - http_code = HttpCodes.SUCCESS if mail_data else HttpCodes.NOT_FOUND, - data = mail_data + status_code = StatusCodes.OK if mails_list else StatusCodes.FAILED, + http_code = HttpCodes.SUCCESS if mails_list else HttpCodes.NOT_FOUND, + data = mails_list, + message = f"{len(mails_list) if mails_list else 0} mail(s) found" ) diff --git a/api/blueprints/mail/oauth_callback.py b/api/blueprints/mail/oauth_callback.py index 8759b46..730cfd2 100644 --- a/api/blueprints/mail/oauth_callback.py +++ b/api/blueprints/mail/oauth_callback.py @@ -62,9 +62,6 @@ from utils_v2.api.async_quart import ( # Common: from shared import constants -# Behaviour Models: -from models.behaviour.mail.oauth import MailOAuthModel - # For asynchronous activities: import asyncio diff --git a/api/blueprints/mail/retrieve.py b/api/blueprints/mail/retrieve.py index e8d8712..59bf04a 100644 --- a/api/blueprints/mail/retrieve.py +++ b/api/blueprints/mail/retrieve.py @@ -10,7 +10,7 @@ OBJECTIVE: - To enlist multiple e-mails for a given user at a time. + To get one full mail for any given user for any given account. REFERENCES: @@ -48,6 +48,7 @@ from utils_v2.database.async_mongo_v2 import AsyncMongo from utils_v2.api.codes import StatusCodes, HttpCodes from utils_v2.api.response import ResponseModel from utils_v2.api.async_quart import ( + make_ordered_json, set_api_version, read_input, get_session_info, @@ -68,8 +69,7 @@ from utils_v2.goog.models.data.auth_tokens import GoogleAuthTokens from shared import constants # Data Models: -from models.data.mail.sync import MailSyncRequestHeaders, MailSyncRequestData -from models.data.mail.sync import MailSyncOneResult, MailSyncManyResults +from models.data.mail.retrieve import MailGetRequestHeaders, MailGetRequestData # To work with datatypes: from typing import Literal @@ -77,9 +77,6 @@ from typing import Literal # For asynchronous activities: import asyncio -# To work with LLMs: -from langchain_openai import ChatOpenAI - # To work with date and time: import datetime @@ -92,7 +89,7 @@ import datetime # Related to Quart: -mail_sync_bp = Blueprint("mail_sync", __name__) +mail_retrieve_bp = Blueprint("mail_retrieve", __name__) # ***************************************************************************************************************** @@ -112,7 +109,7 @@ mail_sync_bp = Blueprint("mail_sync", __name__) # ***************************************************************************************************************** -@mail_sync_bp.record_once +@mail_retrieve_bp.record_once def init(blueprint_setup_state): # This gets called when the blueprint is registered. @@ -123,41 +120,7 @@ def init(blueprint_setup_state): # --------------------------------------------------------------------------------------------------------------------- -async def sync_mails( - mongo_conn: AsyncMongo, - llm: ChatOpenAI, - inbound_headers: dict, - inbound_data: MailSyncRequestData -) -> MailSyncManyResults: - - """ - A very simple function, but kept separate so that we get the option to switch between running it in the foreground - and running it in the background. - :param mongo_conn: The instance of the database connector to use to sync the mails. - :param llm: The instance of the LLM to use to summarize the mails. - :param inbound_headers: The headers that came in with the request. - :param inbound_data: The data that came in with the request. - :return: The results of the mail-sync'ing attempt. - """ - - # Try to sync the mails: - return await current_app.mail_sync_model.sync( - session_token = inbound_headers["X-Session-Token"], - mongo_conn = mongo_conn, - account_identifier = inbound_data.accountId, - llm = llm, - force_sync = inbound_data.forceSync, - start_date = inbound_data.startDate, - end_date = inbound_data.endDate, - max_count = inbound_data.maxCount - ) - - -# --------------------------------------------------------------------------------------------------------------------- - - -@mail_sync_bp.route("/sync", methods = ["POST"]) -@mail_sync_bp.route("/sync/", methods = ["POST"]) +@mail_retrieve_bp.route("/get", 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") @@ -165,7 +128,7 @@ async def sync_mails( attr_name = "logs_mongo", project = constants.PROJECT_NAME, log_type = constants.MODULE_NAME, - operation = "mailOAuthUrlReqApi", + operation = "mailGetApi", log_input = True, log_output = True, sensitive_keys = ["sessionToken", "X-Session-Token"] @@ -173,22 +136,19 @@ async def sync_mails( @log_chain_to_mongo(attr_name = "logs_mongo") @should_not_be_under_maintenance(attr_name = "is_under_maintenance") @validate_input( - header_validator = lambda x: MailSyncRequestHeaders(**x).model_dump(), - data_validator = lambda x: MailSyncRequestData(**x) + header_validator = lambda x: MailGetRequestHeaders(**x).model_dump(), + data_validator = lambda x: MailGetRequestData(**x) ) @handle_cancelled_request() -async def sync_mail( - mode: Literal["background", "bg"] = None, - inbound_headers: dict | MailSyncRequestHeaders = None, - inbound_data: dict | MailSyncRequestData = None, +async def get_one_mail( + inbound_headers: dict | MailGetRequestHeaders = None, + inbound_data: dict | MailGetRequestData = None, inbound_files: dict = None, **kwargs ): """ - Use this when the user wants to pull old mails from some mail client (like GMail) and save it to the database for - ready access on the UI. - :param mode: Set it to one of the specified options to make the sync'ing process go to the background. + Use this endpoint when the user wants to fetch one mail. :param inbound_headers: auto-extracted by the decorators. :param inbound_data: auto-extracted by the decorators. :param inbound_files: auto-extracted by the decorators. @@ -203,43 +163,17 @@ async def sync_mail( http_code = HttpCodes.UNAUTHORIZED ) - # Make the variables available in the scope of the current request: - g.inbound_headers = inbound_headers - g.inbound_data = inbound_data - - # If we've been asked to sync the mails in the background: - if mode in ["background", "bg"]: - current_app.add_background_task( - sync_mails, - mongo_conn = current_app.data_mongo, - llm = current_app.llm, - inbound_headers = inbound_headers, - inbound_data = inbound_data - ) - return ResponseModel( - status_code = StatusCodes.OK, - http_code = HttpCodes.ACCEPTED, - message = "your mails are being sync'd in the background" - ) - - # Otherwise we process it right here: - sync_results = await sync_mails( + # Get the mail: + mail_data = await current_app.mail_retrieve_model.get_mail( mongo_conn = current_app.data_mongo, - llm = current_app.llm, - inbound_headers = inbound_headers, - inbound_data = inbound_data + mail_identifier = inbound_data.mailId ) - # Response: + # Done here: return ResponseModel( - status_code = StatusCodes.FAILED if sync_results.failureCount > 0 else StatusCodes.OK, - http_code = HttpCodes.INTERNAL_SERVER_ERROR if sync_results.failureCount > 0 else HttpCodes.SUCCESS, - message = sync_results.message, - data = { - "totalCount": sync_results.totalCount, - "successCount": sync_results.successCount, - "failureCount": sync_results.failureCount - } + status_code = StatusCodes.OK if mail_data else StatusCodes.FAILED, + http_code = HttpCodes.SUCCESS if mail_data else HttpCodes.NOT_FOUND, + data = mail_data ) diff --git a/api/blueprints/mail/sync.py b/api/blueprints/mail/sync.py index b0778a6..04487ea 100644 --- a/api/blueprints/mail/sync.py +++ b/api/blueprints/mail/sync.py @@ -167,7 +167,7 @@ async def sync_mails( attr_name = "logs_mongo", project = constants.PROJECT_NAME, log_type = constants.MODULE_NAME, - operation = "mailOAuthUrlReqApi", + operation = "mailSyncApi", log_input = True, log_output = True, sensitive_keys = ["sessionToken", "X-Session-Token"] diff --git a/api/main.py b/api/main.py index 78aa166..7c11cd7 100644 --- a/api/main.py +++ b/api/main.py @@ -73,6 +73,7 @@ from utils_v2.goog.gmail.gmail_client import AsyncGMailClient # Behaviour Models: from models.behaviour.mail.oauth_v2 import MailOAuthModel from models.behaviour.mail.sync_v2 import MailSyncModel +from models.behaviour.mail.retrieve import MailRetrieveModel # To make REST API calls: import httpx @@ -84,6 +85,8 @@ from icecream import IceCreamDebugger from api.blueprints.mail.oauth_request import mail_oauth_bp from api.blueprints.mail.oauth_callback import mail_callback_bp from api.blueprints.mail.sync import mail_sync_bp +from api.blueprints.mail.list import mail_list_bp +from api.blueprints.mail.retrieve import mail_retrieve_bp from api.blueprints.tech.chat_alerts import tech_chat_alert_bp from api.blueprints.test.callback import test_callback_bp @@ -119,6 +122,8 @@ app = cors(app) app.register_blueprint(mail_oauth_bp, url_prefix = f"/{MODULE_BASE}/mail") app.register_blueprint(mail_callback_bp, url_prefix = f"/{MODULE_BASE}/mail") app.register_blueprint(mail_sync_bp, url_prefix = f"/{MODULE_BASE}/mail") +app.register_blueprint(mail_list_bp, url_prefix = f"/{MODULE_BASE}/mail") +app.register_blueprint(mail_retrieve_bp, url_prefix = f"/{MODULE_BASE}/mail") app.register_blueprint(tech_chat_alert_bp, url_prefix = f"/{MODULE_BASE}/tech/alert") app.register_blueprint(test_callback_bp, url_prefix = f"/{MODULE_BASE}/test") @@ -309,10 +314,18 @@ async def app_startup(**kwargs): debug_prefix = "Mail-Sync | ", debug_only_errors = True ) + current_app.mail_retrieve_model = MailRetrieveModel( + cache = current_app.module_cache, + alert_url = current_app.script_data["alerts"]["url"], + http_client = current_app.http_client, + debug = enable_debugging, + debug_prefix = "Mail-Retr | ", + debug_only_errors = True + ) - # ┏┓ - # ┃ ┏┓┏┓┏┓┏┓┏╋┏┓┏┓┏ - # ┗┛┗┛┛┗┛┗┗ ┗┗┗┛┛ ┛ + # ┏┓ ┓ ┏┓┓• + # ┃ ┏┓┏┓┏┓┏┓┏╋┏┓┏┓┏ ┏┓┏┓┏┫ ┃ ┃┓┏┓┏┓╋┏ + # ┗┛┗┛┛┗┛┗┗ ┗┗┗┛┛ ┛ ┗┻┛┗┗┻ ┗┛┗┗┗ ┛┗┗┛ # Create an instance to handle GMail-related activities: current_app.gmail_client = AsyncGMailClient( @@ -325,6 +338,11 @@ async def app_startup(**kwargs): debug_only_errors = False ) + # ┏┓┳ ┳┳┓ • + # ┣┫┃ ┃┃┃┏┓┏┓┓┏ + # ┛┗┻ ┛ ┗┗┻┗┫┗┗ + # ┛ + # For LLMs: current_app.llm = ChatOpenAI( model = script_cred["openAi"]["model"], diff --git a/models/behaviour/mail/retrieve.py b/models/behaviour/mail/retrieve.py index abaf213..c68f956 100644 --- a/models/behaviour/mail/retrieve.py +++ b/models/behaviour/mail/retrieve.py @@ -108,11 +108,12 @@ class MailRetrieveModel(BaseModel): :return: Either the JSON that describes the mail or None if such a mail does not exist. """ + # Get the data from the database: mail_data = await mongo_conn.find_one( collection = self.MAIL_COLLECTION, filter = {"_id": ObjectId(mail_identifier)}, projection = { - "mailId": "_id", + "_id": True, "serviceType": True, "client": True, "payload.ts": True, @@ -125,13 +126,22 @@ class MailRetrieveModel(BaseModel): "payload.attachments": True, "payload.labels": True, "payload.snippet": True, - "payload.aiSnippet": "payload.aiSnippet.snippet", + "payload.aiSnippet": True, } ) - if mail_data: mail_data["mailId"] = str(mail_data["mailId"]) + + # Format the data: + if mail_data: + mail_data["mailId"] = str(mail_data.pop("_id")) + mail_data["payload"]["ts"] = mail_data["payload"]["ts"].isoformat() + mail_data["payload"]["readTs"] = mail_data["payload"]["readTs"].isoformat() + if ai_snippet := mail_data["payload"].pop("aiSnippet"): + mail_data["payload"]["aiSnippet"] = ai_snippet["snippet"] + + # Done here: return mail_data - async def list_by_account_identifier( + async def list_for_account_identifier( self, mongo_conn: AsyncMongo, account_identifier: str | ObjectId, @@ -139,7 +149,46 @@ class MailRetrieveModel(BaseModel): skip: int = 0 ): - pass + """ + To enlist mails for one account. + :param mongo_conn: The instance of the database connector to use to get the mail's data. + :param account_identifier: The id of the document in the database that holds the tokens to access the account. + :param limit: How many records to fetch. + :param skip: How many initial records to skip. useful for pagination. + :return: Either the JSON that describes the mails or None if something failed. + """ + + # Get the data from the database: + mails_list = await mongo_conn.find_many( + collection = self.MAIL_COLLECTION, + filter = {"accountId": ObjectId(account_identifier)}, + projection = { + "_id": True, + "serviceType": True, + "client": True, + "payload.ts": True, + "payload.readTs": True, + "payload.from": True, + "payload.labels": True, + "payload.snippet": True, + "payload.aiSnippet": True, + }, + limit = limit, + skip = skip, + sort = {"payload.ts": -1} + ) + + # Format the data: + if mails_list: + for mail_data in mails_list: + mail_data["mailId"] = str(mail_data.pop("_id")) + mail_data["payload"]["ts"] = mail_data["payload"]["ts"].isoformat() + mail_data["payload"]["readTs"] = mail_data["payload"]["readTs"].isoformat() + if ai_snippet := mail_data["payload"].pop("aiSnippet"): + mail_data["payload"]["aiSnippet"] = ai_snippet["snippet"] + + # Done here: + return mails_list # ***************************************************************************************************************** diff --git a/models/data/mail/list.py b/models/data/mail/list.py index d4e6abb..e7b9805 100644 --- a/models/data/mail/list.py +++ b/models/data/mail/list.py @@ -75,7 +75,7 @@ REGEX_SESSION_TOKEN = r"^[a-f0-9]{8}-[a-f0-9]{4}-[1-5][a-f0-9]{3}-[89ab][a-f0-9] # ***************************************************************************************************************** -class MailGetRequestHeaders(BaseModel): +class MailListRequestHeaders(BaseModel): sessionToken: str = Field( description = "the session token of the user who is requesting the service", @@ -99,10 +99,24 @@ class MailGetRequestHeaders(BaseModel): # --------------------------------------------------------------------------------------------------------------------- -class MailGetRequestData(BaseModel): +class MailListByAccountIdRequestData(BaseModel): - mailId: str = Field( - description = "the mail identifier (Mongo ObjectId) of the document that holds the mail", + accountId: str = Field( + description = "the account identifier (Mongo ObjectId) granted by 'MailOAuthModel.get_account_identifier'", + 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", + default = 0, frozen = True ) diff --git a/models/data/mail/retrieve.py b/models/data/mail/retrieve.py index ee3965c..d4e6abb 100644 --- a/models/data/mail/retrieve.py +++ b/models/data/mail/retrieve.py @@ -6,11 +6,11 @@ DATE: - Monday, 2nd Dec., 2024. + Tuesday, 3rd Dec., 2024. OBJECTIVE: - To provide the structure for the request that will come in to sync the mails of a particular user. + To provide a structure to query the full payload of an email. REFERENCES: @@ -75,7 +75,7 @@ REGEX_SESSION_TOKEN = r"^[a-f0-9]{8}-[a-f0-9]{4}-[1-5][a-f0-9]{3}-[89ab][a-f0-9] # ***************************************************************************************************************** -class MailSyncRequestHeaders(BaseModel): +class MailGetRequestHeaders(BaseModel): sessionToken: str = Field( description = "the session token of the user who is requesting the service", @@ -99,114 +99,13 @@ class MailSyncRequestHeaders(BaseModel): # --------------------------------------------------------------------------------------------------------------------- -class MailSyncRequestData(BaseModel): +class MailGetRequestData(BaseModel): - accountId: str = Field( - description = "the account identifier (Mongo ObjectId) granted by 'MailOAuthModel.get_account_identifier'", + mailId: str = Field( + description = "the mail identifier (Mongo ObjectId) of the document that holds the mail", frozen = True ) - maxCount: int = Field( - description = "the max. no. of e-mails to sync at a given time", - default = 100, - ge = 1, - le = 100, - frozen = True - ) - - startDate: PastDatetime = Field( - description = "the starting date from which the user wants to sync their mail", - default_factory = lambda: date_time.get_current_utc_date_time() - datetime.timedelta(days = 1), - frozen = True - ) - - endDate: PastDatetime = Field( - description = "the ending date till which the user wants to sync their mail", - default_factory = lambda: date_time.get_current_utc_date_time() - datetime.timedelta(seconds = 1), - frozen = True - ) - - forceSync: bool = Field( - description = "use this to forcefully re-sync mails when you need to overwrite existing data in mongodb", - default = False - ) - - # ┏┓ ┏• - # ┃ ┏┓┏┓╋┓┏┓ - # ┗┛┗┛┛┗┛┗┗┫ - # ┛ - - class Config: - extra = "forbid" - - # ┓┏ ┓• ┓ • - # ┃┃┏┓┃┓┏┫┏┓╋┓┏┓┏┓ - # ┗┛┗┻┗┗┗┻┗┻┗┗┗┛┛┗ - - @field_validator("startDate", "endDate", mode = "before") - def to_datetime(cls, value): - if not isinstance(value, datetime.datetime): - value = date_time.parse_date_time( - input_value = value, - timezone = date_time.TIMEZONE_UTC - ) - return value - - -# --------------------------------------------------------------------------------------------------------------------- - - -class MailSyncOneResult(BaseModel): - - success: bool = Field( - description = "whether, or not, the mail was successfully sync'd", - default = False - ) - - message: str | None = Field( - description = "a brief message to summarize the result of the process", - default = None - ) - - mailMessage: dict | None = Field( - description = "the actual data of the mail; can be null in a successful process if the mail is already sync'd", - default = None - ) - - # ┏┓ ┏• - # ┃ ┏┓┏┓╋┓┏┓ - # ┗┛┗┛┛┗┛┗┗┫ - # ┛ - - class Config: - extra = "forbid" - - -# --------------------------------------------------------------------------------------------------------------------- - - -class MailSyncManyResults(BaseModel): - - totalCount: int = Field( - description = "the total no. of mails that were to be sync'd", - default = 0 - ) - - successCount: int = Field( - description = "the no. of mails that were successfully sync'd", - default = 0 - ) - - failureCount: int = Field( - description = "the no. of mails that were successfully sync'd", - default = 0 - ) - - message: str = Field( - description = "a brief message to summarize the results of the process", - default = None - ) - # ┏┓ ┏• # ┃ ┏┓┏┓╋┓┏┓ # ┗┛┗┛┛┗┛┗┗┫