import datetime from concurrent.futures import ThreadPoolExecutor, wait import logging from typing import Any, Dict import boto3 from django.conf import settings from django.core.cache import cache from django_redis.cache import RedisCache from django.urls import reverse from sorl.thumbnail.shortcuts import get_thumbnail, delete logger = logging.getLogger(__name__) cloudfront = boto3.client("cloudfront") __executor = None def get_executor() -> ThreadPoolExecutor: global __executor if __executor is None: __executor = ThreadPoolExecutor() return __executor class AsyncRedisCache(RedisCache): def __init__(self, server: str, params: Dict[str, Any]): super().__init__(server=server, params=params) def set(self, *args, **kwargs): return self.apply_async(super().set, *args, **kwargs) def set_many(self, *args, **kwargs): return self.apply_async(super().set_many, *args, **kwargs) def apply_async(self, func, *args, **kwargs): if settings.DJANGO_REDIS_IS_ASYNC: get_executor().submit(func, *args, **kwargs) else: func(*args, **kwargs) class ServiceCacheMixin: CACHE_TTL = settings.CACHE_TTL def cache_get(self, keys): return cache.get(self.cache_key([self.__class__.__name__, *keys])) def cache_set(self, keys, value, ttl=None): cache.set( self.cache_key([self.__class__.__name__, *keys]), value, ttl or self.CACHE_TTL, ) def cache_set_many(self, data, ttl=None, dry_run=False): args = [ { self.cache_key((self.__class__.__name__, *k)): v for k, v in data.items() }, ttl or self.CACHE_TTL, ] if dry_run: return args cache.set_many(*args) def cache_delete(self, keys): return cache.delete(self.cache_key([self.__class__.__name__, *keys])) @staticmethod def cache_key(keys): keys = [str(key) for key in keys] return "-".join(keys) def invalidate_cloudfront_by_id(item_id): if not settings.CLOUDFRONT_DISTRIBUTION_ID or not item_id: return items = ( reverse("artist_by_gras_participant_id", args=[f"{item_id}*"]), reverse("artist_by_gras_participant_ids", args=[f"*{item_id}*"]), reverse("label_by_label_id", args=[f"{item_id}*"]), reverse("label_by_label_id", args=[f"*{item_id}*"]), ) now = datetime.datetime.now() caller_ref = f"{item_id} - {now:%Y/%m/%d %H:%M}" distribution_id = settings.CLOUDFRONT_DISTRIBUTION_ID try: cloudfront.create_invalidation( DistributionId=distribution_id, InvalidationBatch={ "Paths": {"Quantity": len(items), "Items": items}, "CallerReference": caller_ref, }, ) except Exception as e: logger.exception( f"Error on CloudFront cache invalidation: {item_id}, {e}" ) else: logger.info(f"CloudFront cache invalidation: {item_id}") def invalidate_cloudfront_by_ids(item_ids): wait( get_executor().submit(invalidate_cloudfront_by_id, item_id) for item_id in item_ids ) def create_image_thumbnail(image, resolution): try: thumbnail = get_thumbnail( image, resolution, upscale=False, crop=False, ) logger.info( f"Created thumbnail: {image} -> {resolution}: {thumbnail.url}" # noqa ) except Exception as e: logger.error(f"Unable to get thumbnail: {image.name} - {e}") raise e def create_image_thumbnails(image, resolutions, thread=True): if thread: wait( get_executor().submit(create_image_thumbnail, image, resolution) for resolution in resolutions ) else: for resolution in resolutions: create_image_thumbnail(image, resolution) def delete_image_thumbnails(image): get_executor().submit(delete, image)