(20241225) Small changes in Async Mongo.

This commit is contained in:
2024-12-25 19:47:47 +05:30
parent 1fd287900d
commit 7607a427a6
6 changed files with 587 additions and 113 deletions
+401
View File
@@ -0,0 +1,401 @@
"""
AUTHOR:
Khushal P Soonderji
DATE:
Created: Wednesday, 18th Sept., 2024
Updated: Wednesday, 25th Dec., 2024
OBJECTIVE:
To be able to fetch logs for rapid issue resolution.
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
# 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,
log_request_to_mongo,
should_not_be_under_maintenance,
only_whitelisted_ips,
limit_rate,
validate_input,
handle_cancelled_request
)
# Models:
from utils_v2.api.log import APILogModel
from models.logs.api import LogChainRequestData, LogsByFilterRequestData
# For asynchronous activities:
import asyncio
# *****************************************************************************************************************
# ***** ****
# *** MACROS / ONE-TIME INIT ***
# ***** ****
# *****************************************************************************************************************
# Related to Quart:
logs_bp = Blueprint("int_logs", __name__)
# *****************************************************************************************************************
# ***** ****
# *** VARIABLES ***
# ***** ****
# *****************************************************************************************************************
# --- Nothing Yet
# *****************************************************************************************************************
# ***** ****
# *** FUNCTIONS ***
# ***** ****
# *****************************************************************************************************************
@logs_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
# ---------------------------------------------------------------------------------------------------------------------
@logs_bp.route("/get/id/<log_id>", methods = ["POST", "GET"])
@set_api_version(api_version = "2.1.0")
@should_not_be_under_maintenance(attr_name = "is_under_maintenance")
@only_whitelisted_ips(attr_name = "whitelisted_ips")
@handle_cancelled_request()
async def get_log(
log_id,
**kwargs
):
"""
To get the log from its log id.
:param log_id: An identifier (string) for the log to fetch.
"""
fetched_log = await current_app.mongo.find_one(
collection = "logs",
filter = {"logId": log_id},
projection = {"_id": False}
)
if not fetched_log: return ResponseModel(
status_code = StatusCodes.FAILED,
http_code = HttpCodes.NOT_FOUND,
message = "No such log."
)
else: return ResponseModel(
status_code = StatusCodes.OK,
message = f"Log found.",
data = {"total": 1, "fetched": 1, "logs": fetched_log}
)
# ---------------------------------------------------------------------------------------------------------------------
@logs_bp.route("/get/exception/id/<log_id>", methods = ["POST", "GET"])
@set_api_version(api_version = "2.1.0")
@should_not_be_under_maintenance(attr_name = "is_under_maintenance")
@only_whitelisted_ips(attr_name = "whitelisted_ips")
async def get_exception_from_log(
log_id,
**kwargs
):
"""
To get the log's exception from its log id.
:param log_id: An identifier (string) for the log to fetch.
"""
fetched_log = await current_app.mongo.find_one(
collection = "logs",
filter = {"logId": log_id},
projection = {
"_id": False,
"logId": True,
"log": True,
"operation": True,
"ts": True,
"exception": True
}
)
if not fetched_log: return ResponseModel(
status_code = StatusCodes.FAILED,
http_code = HttpCodes.NOT_FOUND,
message = "No such log."
)
else: return ResponseModel(
status_code = StatusCodes.OK,
message = f"Log found.",
data = {"total": 1, "fetched": 1, "logs": fetched_log}
)
# ---------------------------------------------------------------------------------------------------------------------
@logs_bp.route("/get/chain/<log_chain>", methods = ["POST", "GET"])
@set_api_version(api_version = "2.1.0")
@read_input(sanitize_headers = False, sanitize_data = False)
@should_not_be_under_maintenance(attr_name = "is_under_maintenance")
@only_whitelisted_ips(attr_name = "whitelisted_ips")
@validate_input(data_validator = lambda x: LogChainRequestData(**x))
@handle_cancelled_request()
async def get_log_chain(
log_chain,
inbound_headers: dict = None,
inbound_data: dict | LogChainRequestData = None,
inbound_files: dict = None,
**kwargs
):
"""
To get the series of logs from its chain identifier.
:param log_chain: An identifier (string) for the log chain to fetch.
: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.
"""
total_count = await current_app.mongo.count(
collection = "logs",
filter = {"logChain": log_chain}
)
fetched_logs = await current_app.mongo.find_many(
collection = "logs",
filter = {"logChain": log_chain},
projection = inbound_data.projection,
sort = inbound_data.sort,
limit = inbound_data.limit,
skip = inbound_data.skip
)
fetched_count = len(fetched_logs)
if not fetched_logs: return ResponseModel(
status_code = StatusCodes.FAILED,
http_code = HttpCodes.NOT_FOUND,
message = "No such log chain."
)
else: return ResponseModel(
status_code = StatusCodes.OK,
message = f"{fetched_count} log(s) fetched",
data = {"total": total_count, "fetched": fetched_count, "logs": fetched_logs}
)
# ---------------------------------------------------------------------------------------------------------------------
@logs_bp.route("/get/exception/chain/<log_chain>", methods = ["POST", "GET"])
@set_api_version(api_version = "2.1.0")
@read_input(sanitize_headers = False, sanitize_data = False)
@should_not_be_under_maintenance(attr_name = "is_under_maintenance")
@only_whitelisted_ips(attr_name = "whitelisted_ips")
@validate_input(data_validator = lambda x: LogChainRequestData(**x))
@handle_cancelled_request()
async def get_exceptions_from_log_chain(
log_chain,
inbound_headers: dict = None,
inbound_data: dict | LogChainRequestData = None,
inbound_files: dict = None,
**kwargs
):
"""
To get the series of log exceptions from its chain identifier.
:param log_chain: An identifier (string) for the log chain to fetch.
: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.
"""
total_count = await current_app.mongo.count(
collection = "logs",
filter = {"logChain": log_chain}
)
fetched_logs = await current_app.mongo.find_many(
collection = "logs",
filter = {"logChain": log_chain},
projection = {
"_id": False,
"log": True,
"operation": True,
"ts": True,
"exception": True
},
sort = inbound_data.sort,
limit = inbound_data.limit,
skip = inbound_data.skip
)
fetched_count = len(fetched_logs)
if not fetched_logs: return ResponseModel(
status_code = StatusCodes.FAILED,
http_code = HttpCodes.NOT_FOUND,
message = "No such log chain."
)
else: return ResponseModel(
status_code = StatusCodes.OK,
message = f"{fetched_count} log(s) fetched.",
data = {"total": total_count, "fetched": fetched_count, "logs": fetched_logs}
)
# ---------------------------------------------------------------------------------------------------------------------
@logs_bp.route("/get/filter", methods = ["POST", "GET"])
@set_api_version(api_version = "2.1.0")
@read_input(sanitize_headers = False, sanitize_data = False)
@should_not_be_under_maintenance(attr_name = "is_under_maintenance")
@only_whitelisted_ips(attr_name = "whitelisted_ips")
@validate_input(data_validator = lambda x: LogsByFilterRequestData(**x))
@handle_cancelled_request()
async def get_logs_by_filter(
inbound_headers: dict = None,
inbound_data: dict | LogsByFilterRequestData = None,
inbound_files: dict = None,
**kwargs
):
"""
To fetch logs by custom filters:
: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.
"""
total_count = await current_app.mongo.count(
collection = "logs",
filter = inbound_data.filter
)
fetched_logs = await current_app.mongo.find_many(
collection = "logs",
filter = inbound_data.filter,
projection = inbound_data.projection,
sort = inbound_data.sort,
limit = inbound_data.limit,
skip = inbound_data.skip
)
fetched_count = len(fetched_logs)
if not fetched_logs: return ResponseModel(
status_code = StatusCodes.FAILED,
http_code = HttpCodes.NOT_FOUND,
message = "No matching logs."
)
else: return ResponseModel(
status_code = StatusCodes.OK,
message = f"{len(fetched_logs)} / {total_count} log(s) fetched.",
data = {"total": total_count, "fetched": fetched_count, "logs": fetched_logs}
)
# ---------------------------------------------------------------------------------------------------------------------
@logs_bp.route("/set/api", methods = ["POST"])
@set_api_version(api_version = "2.1.0")
@read_input(sanitize_headers = True, sanitize_data = True)
@should_not_be_under_maintenance(attr_name = "is_under_maintenance")
@only_whitelisted_ips(attr_name = "whitelisted_ips")
@validate_input(data_validator = lambda x: APILogModel(**x))
async def set_api_log(
inbound_headers: dict = None,
inbound_data: dict | APILogModel = None,
inbound_files: dict = None,
**kwargs
):
"""
To set logs from internal whitelisted IPs.
: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.
"""
inserted_id = await current_app.mongo.insert_one(
collection = "logs",
document = inbound_data.model_dump()
)
return ResponseModel(
status_code = StatusCodes.OK if inserted_id else StatusCodes.FAILED,
http_code = HttpCodes.SUCCESS if inserted_id else HttpCodes.BAD_REQUEST
)
# *****************************************************************************************************************
# ***** ****
# *** MAIN PROGRAM ***
# ***** ****
# *****************************************************************************************************************
if __name__ == "__main__":
pass
+184 -111
View File
@@ -10,7 +10,7 @@
OBJECTIVE: OBJECTIVE:
To provide a structure to work with requests surrounding Cred and Data handling. To provide a structure to work with enlisting servers and maintaining their status.
REFERENCES: REFERENCES:
@@ -36,13 +36,17 @@ sys.path.append(".")
sys.path.append("..") sys.path.append("..")
# For making data behaviour_models: # For making data behaviour_models:
from pydantic import BaseModel, Field, field_validator, PastDatetime, model_validator from pydantic import BaseModel, Field, field_validator, PastDatetime, model_validator, AwareDatetime, computed_field
from typing import Optional, Literal, Union from typing import Optional, Literal, Union
# My utils: # My utils:
from utils_v2.string import json
from utils_v2.string import regex from utils_v2.string import regex
from utils_v2.date_time import date_time from utils_v2.date_time import date_time
# To work with MongoDB:
from bson import ObjectId
# To work with date and time: # To work with date and time:
import datetime import datetime
@@ -74,134 +78,178 @@ import datetime
# ***************************************************************************************************************** # *****************************************************************************************************************
class CredAndDataSetRequestHeaders(BaseModel): class CoreServerInfoModel(BaseModel):
scriptId: str = Field( serverId: ObjectId | None = Field(
description = "the id of the script for whom you are setting cred/data", description = "the id of the document in mongodb that holds the info. about this server",
default = None,
frozen = True, frozen = True,
alias = "X-Script-Id" exclude = True,
alias = "_id"
) )
scriptDescription: str = Field( hostname: str = Field(
description = "a short description of the script and what it does", description = "the hostname to identify the server",
frozen = True, frozen = True
alias = "X-Script-Desc"
) )
os: str = Field(
description = "the os the server is running",
frozen = True
)
cpu: str = Field(
description = "the cpu that the server has in it",
frozen = True
)
pid: str | int = Field(
description = "the process id that registered the details of this server",
frozen = True
)
ppid: str | int = Field(
description = "the parent process id that registered the details of this server",
frozen = True
)
ipAddr: str | None = Field(
description = "the ip address of the server",
frozen = True
)
portNo: int | None = Field(
description = "the port no. that this service is running on",
default = None,
frozen = True
)
project: str = Field(
description = "the name of the project that this server is running",
frozen = True
)
service: str = Field(
description = "the service in the said project that this server is running",
frozen = True
)
description: str = Field(
description = "a brief description about this service",
max_length = 350,
frozen = True
)
healthCheckUrl: str = Field(
description = "a get request will be sent to this url to see if the server is up",
frozen = True
)
healthCheckInterval: int = Field(
description = "the seconds after which to check for the service being up",
ge = 30,
frozen = True
)
healthAlertUrl: str = Field(
description = "a get request will be sent to this url when the server goes offline",
frozen = True
)
online: bool = Field(
description = "whether this service is online or not",
default = False,
frozen = False
)
batchId: str | None = Field(
description = "set some value here when checking the status of this server, set it to null once done",
default = None,
frozen = False
)
batchTs: AwareDatetime | None = Field(
description = "set the time (utc) when checking the status of this server, set it to null once done",
default = None,
frozen = False
)
firstRegTs: AwareDatetime = Field(
description = "the first time (utc) this server was registered in the database",
default_factory = lambda: date_time.get_current_utc_date_time(as_string = False),
frozen = True
)
lastRegTs: AwareDatetime = Field(
description = "the last time (utc) this server was registered in the database",
default_factory = lambda: date_time.get_current_utc_date_time(as_string=False),
frozen = True
)
lastCheckTs: AwareDatetime | None = Field(
description = "the last time (utc) this server was checked for being online",
default = None,
frozen = True
)
@computed_field
def checkAfterTs(self) -> datetime.datetime:
if self.lastCheckTs is None: return date_time.get_current_utc_date_time(as_string = False)
else: return self.lastCheckTs + datetime.timedelta(seconds = self.healthCheckInterval)
# ┏┓ ┏• # ┏┓ ┏•
# ┃ ┏┓┏┓╋┓┏┓ # ┃ ┏┓┏┓╋┓┏┓
# ┗┛┗┛┛┗┛┗┗┫ # ┗┛┗┛┛┗┛┗┗┫
# ┛ # ┛
class Config: class Config:
extra = "allow" extra = "ignore"
arbitrary_types_allowed = True
def model_dump(self, *args, **kwargs):
return super().model_dump(*args, by_alias = True, **kwargs)
# ---------------------------------------------------------------------------------------------------------------------
class CredAndDataGetRequestHeaders(BaseModel):
scriptId: str = Field(
description = "the id of the script for whom you are getting cred/data",
frozen = True,
alias = "X-Script-Id"
)
# ┏┓ ┏•
# ┃ ┏┓┏┓╋┓┏┓
# ┗┛┗┛┛┗┛┗┗┫
# ┛
class Config:
extra = "allow"
def model_dump(self, *args, **kwargs):
return super().model_dump(*args, by_alias = True, **kwargs)
# ---------------------------------------------------------------------------------------------------------------------
class CredAndDataUpdateRequestHeaders(BaseModel):
scriptId: str = Field(
description = "the id of the script for whom you are updating cred/data",
frozen = True,
alias = "X-Script-Id"
)
# ┏┓ ┏•
# ┃ ┏┓┏┓╋┓┏┓
# ┗┛┗┛┛┗┛┗┗┫
# ┛
class Config:
extra = "allow"
def model_dump(self, *args, **kwargs):
return super().model_dump(*args, by_alias = True, **kwargs)
# ---------------------------------------------------------------------------------------------------------------------
class CredAndDataUpdateRequestData(BaseModel):
unsetJson: dict = Field(
description = "the items to unset; this is performed first",
frozen = True,
alias = "unset"
)
setJson: dict = Field(
description = "the items to set; this is performed after the 'unset' operation",
frozen = True,
alias = "set"
)
# ┏┓ ┏•
# ┃ ┏┓┏┓╋┓┏┓
# ┗┛┗┛┛┗┛┗┗┫
# ┛
class Config:
extra = "forbid"
# ┓┏ ┓• ┓ • # ┓┏ ┓• ┓ •
# ┃┃┏┓┃┓┏┫┏┓╋┓┏┓┏┓ # ┃┃┏┓┃┓┏┫┏┓╋┓┏┓┏┓
# ┗┛┗┻┗┗┗┻┗┻┗┗┗┛┛┗ # ┗┛┗┻┗┗┗┻┗┻┗┗┗┛┛┗
@field_validator("unsetJson", "setJson", mode = "before") @staticmethod
def ensure_non_null(cls, value): def parse_date_time(value):
if value is None: value = {}
# If a null value was given,
# we can't do anything:
if not value: value = None
# If the input is a string:
if isinstance(value, str):
value = value.strip()
value = date_time.parse_date_time(
input_value = value,
timezone = date_time.TIMEZONE_UTC,
date_formats = [
"%Y-%m-%d",
"%Y-%m-%d %M:%H:%S",
"%Y%m%d",
"%Y%m%d %M:%H:%S",
]
)
# If the input is already a date-time object,
# we just normalize the timestamp:
if isinstance(value, datetime.datetime):
value = date_time.to_timezone(
value,
timezone = date_time.TIMEZONE_UTC
)
# Done here:
return value return value
@field_validator(
# --------------------------------------------------------------------------------------------------------------------- "batchTs",
"firstRegTs", "lastRegTs",
"lastCheckTs",
class CredAndDataDeleteRequestHeaders(BaseModel): mode = "before"
scriptId: str = Field(
description = "the id of the script for whom you are deleting cred/data",
frozen = True,
alias = "X-Script-Id"
) )
def parse_given_date_time(cls, value):
# ┏┓ ┏• return cls.parse_date_time(value)
# ┃ ┏┓┏┓╋┓┏┓
# ┗┛┗┛┛┗┛┗┗┫
# ┛
class Config:
extra = "allow"
def model_dump(self, *args, **kwargs):
return super().model_dump(*args, by_alias = True, **kwargs)
# ***************************************************************************************************************** # *****************************************************************************************************************
@@ -213,4 +261,29 @@ class CredAndDataDeleteRequestHeaders(BaseModel):
if __name__ == "__main__": if __name__ == "__main__":
pass test_json = {
"hostname": "kbprod",
"os": "Ubuntu 22.04.5 LTS",
"cpu": "x86_64 (x86_64)",
"pid": 456,
"ppid": 123,
"ipAddr": "127.0.0.1",
"portNo": 5000,
"project": "Bicree",
"service": "user",
"description": "This is Bicree's user module.",
"healthCheckUrl": "https://v2.api.bicree.com/user/metrics/memory",
"healthCheckInterval": 30,
"healthAlertUrl": (
"https://nexcom.ditscentre.in/wtt/webhook/telegram/out"
"?appKey=9999999"
"&chatId=1275560043"
"&message=%F0%9F%9A%A8%20Bicree%27s%20user%20module%20is%20down%21%20Please%20check%20it%20urgently%21"
),
"online": True,
"batchId": None,
"batchTs": None
}
test_model = CoreServerInfoModel(**test_json)
print("TEST MODEL:", json.to_string(test_model.model_dump(), default = str))
+2 -2
View File
@@ -658,7 +658,7 @@ class AsyncMongo(AsyncMongoBase):
:param upsert: If you want to insert if the document doesn't already exist. :param upsert: If you want to insert if the document doesn't already exist.
:param session: The session if you need to do this in a transaction. :param session: The session if you need to do this in a transaction.
:param raise_exception: Whether, or not, you want to raise an exception when something fails. :param raise_exception: Whether, or not, you want to raise an exception when something fails.
:return: True or False based on the success of the operation. :return: The no. of records that were affected.
""" """
# Ensure you are connected: # Ensure you are connected:
@@ -675,7 +675,7 @@ class AsyncMongo(AsyncMongoBase):
upsert = upsert, upsert = upsert,
session = session session = session
) )
update_count = response.modified_count + response.upserted_count update_count = response.modified_count
# When something goes wrong: # When something goes wrong:
except Exception as exception: except Exception as exception: