from datetime import date, timedelta from loguru import logger from ..config import QUERY_PATHS from ..db import get_rdb from ..db.queries import QueryLoader def _get_dates(rdb, ql: QueryLoader) -> tuple[date | None, date | None]: data = rdb.query(ql.load( QUERY_PATHS.get_dates, main="prod_eu_analytics.summer_playlists_stats_main", source_1="spotify.fact_streams", source_2="spotify.fact_streams", )) last_processed = data["last_processed"][0] last_available = data["last_available"][0] if last_processed is None or last_available is None: return last_processed, last_available start_date = last_processed + timedelta(days=1) return start_date, last_available def run_update() -> str | None: logger.info("Starting stats update") rdb = get_rdb() ql = QueryLoader() start_date, end_date = _get_dates(rdb, ql) if start_date is None or end_date is None: logger.warning("Could not determine date range (empty table or no source data)") return None if start_date > end_date: logger.info("No new data") return None date_list = [start_date + timedelta(days=x) for x in range((end_date - start_date).days + 1)] logger.info(f"Updating {len(date_list)} day(s): {date_list[0]} to {date_list[-1]}") for d in date_list: ds = str(d) logger.info(f"Processing {d}") status = rdb.execute(ql.load(QUERY_PATHS.append_stats, date=ds)) if not status: logger.error(f"Failed to append stats for {d}") return None logger.info("Rebuilding stats hyper table") status = rdb.execute(ql.load(QUERY_PATHS.make_stats_hyper)) if not status: logger.error("Failed to rebuild stats hyper table") return None logger.success(f"Stats updated to {date_list[-1]}") return date_list[-1].strftime("%Y-%m-%d")