from datetime import datetime import pymongo from vendor_image_caching import config from vendor_image_caching.cache.base import BaseClient from vendor_image_caching.constants import CacheField, MONGO_URI_TEMPLATE_LOCAL, MONGO_URI_TEMPLATE_NON_LOCAL from vendor_image_caching.utils import compress, decompress class Client(BaseClient): def __init__(self): self.mongo_client: pymongo.MongoClient = self._get_mongo_client() @staticmethod def _get_mongo_client() -> pymongo.MongoClient: """Get mongo client. AWS DocDB reference Download CA wget https://s3.amazonaws.com/rds-downloads/rds-combined-ca-bundle.pem mongodb://{user}:{password}@{host}:{port}/{db}?ssl=true&ssl_ca_certs=rds-combined-ca-bundle.pem&replicaSet=rs0 &readPreference=secondaryPreferred&retryWrites=false Returns: Pymongo client. """ mongo_uri_template = MONGO_URI_TEMPLATE_LOCAL if config.IS_LOCAL else MONGO_URI_TEMPLATE_NON_LOCAL mongo_uri = mongo_uri_template.format( user=config.MONGODB_USER, password=config.MONGODB_PASSWORD, host=config.MONGODB_HOST, port=config.MONGODB_PORT, db=config.MONGODB_DATABASE, auth=config.MONGODB_AUTH_MECHANISM, ) global mongo_client mongo_client = pymongo.MongoClient(mongo_uri) return mongo_client @staticmethod def _get_filter_condition(record_id: str, created_at: datetime, check_timeout: bool = False) -> dict: """Get MongoDB filter condition . Args: record_id: Record ID. created_at: Record created_at timestamp. Returns: Cached data or None. """ return { CacheField.RECORD_ID: record_id, CacheField.CREATED_AT: created_at, CacheField.IMAGES_SAVED: None } def get_document(self, collection_name: str, record_id: str, created_at: datetime) -> dict or None: """Get data from MongoDB. Args: collection_name: Collection name. record_id: Record ID. created_at: Record created_at timestamp. Returns: Cached data or None. """ result = ( self.mongo_client[config.MONGODB_DATABASE][collection_name] .find_one( self._get_filter_condition(record_id, created_at), { CacheField.CREATED_AT: 0, CacheField.OBJECT_ID: 0, CacheField.IMAGES_SAVED: 0, CacheField.RECORD_ID: 0 }, ) ) return decompress(result[CacheField.COMPRESSION], result[CacheField.DATA]) if result else None def save_document(self, collection_name: str, record_id: str, created_at: datetime, data: dict): """Save data to MongoDB. Args: collection_name: Collection name. record_id: Record ID. created_at: Records created_at datetime. data: Data with new image URLs. """ self.mongo_client[config.MONGODB_DATABASE][collection_name].update_one( self._get_filter_condition(record_id, created_at), { "$set": { CacheField.DATA: compress(config.MONGODB_COMPRESSION_LIB, data), CacheField.IMAGES_SAVED: True, }, }, upsert=False )