import signal import time from contextlib import AbstractContextManager from functools import wraps from typing import Callable from redis.lock import Lock as RedisLock from src.config import Config from src.constants import DBType def get_time_str(date_time) -> str: return time.strftime("%Y:%d:%m %H:%M:%S", time.gmtime(date_time)) def timer(logger=None): """Log execution time and add to the result.""" def decorator(f): @wraps(f) def wrapper(*args, **kwargs): time_start = time.time() result = f(*args, **kwargs) exec_time = time.time() - time_start if logger: logger.info(f"{f.__name__} finished in {exec_time} for params: {args}; {kwargs}.") return result, exec_time return wrapper return decorator def time_logger(logger, app_name: str, **_kwargs) -> Callable: """Decorator to log execution time.""" extra_data = " | ".join([f"{k} - {v}" for k, v in _kwargs.items()]) def decorator(f): @wraps(f) def wrapper(*args, **kwargs): time_start = time.time() logger.info(f"{app_name} started at {get_time_str(time_start)} {extra_data}") result = f(*args, **kwargs) time_finish = time.time() logger.info( f"{app_name} finished at {get_time_str(time_finish)} {extra_data} " f"| total execution time in sec - {time_finish - time_start}" ) return result return wrapper return decorator class ControlledLock(AbstractContextManager): """Context manager wrapper for redis lock with extra 'enabled' parameter, non-blocking check and releasing on SIGTERM.""" lock: RedisLock = None def __init__(self, redis_client, key: str, ttl: int, enabled: bool): self.redis_client = redis_client self.lock_key = key self.lock_ttl = ttl self.lock_enabled = enabled self.locked = False if enabled: self.lock = self.redis_client.lock(self.lock_key, timeout=self.lock_ttl) # for deployment: we should release distributed lock on service shutting down signal.signal(signal.SIGTERM, self._sigterm_handler) def __enter__(self): if not self.lock or not self.lock.acquire(blocking=False): return self.locked = True return self.lock def __exit__(self, exc_type, exc_val, exc_tb): self._release() def _sigterm_handler(self, signum, frame): self._release() def _release(self): if self.locked: self.lock.release() def lock(logger, redis_client, config: Config, sentry=None): """Decorator for locks managing and exceptions handling.""" def decorator(f): @wraps(f) def wrapper(*args, **kwargs): with ControlledLock( redis_client, f"{config.APP_INSTANCE_KEY}/lock", config.REDIS_LOCK_TTL, config.REDIS_LOCK_ENABLED ) as lock: try: if config.REDIS_LOCK_ENABLED and not lock: logger.warning(f"{config.APP_INSTANCE_KEY}: can not get lock.") return return f(*args, **kwargs) except Exception as ex: sentry and sentry.capture_exception(ex) logger.error(str(ex)) raise return wrapper return decorator huge_table_model_paths = { "apollo_main_db.apple.models.ApplePlaylistTracklistHistoryReduced2", "apollo_main_db.spotify.models.SpotifyPlaylistTrackListHistory2Reduced2", } def is_huge_table(model_path: str) -> bool: # naive temporary realization, set uo it with env variable instead return model_path in huge_table_model_paths def get_db_type(db_type: str): return DBType[db_type]