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_agg_main", source_1="spotify.agg_playlist_history", source_2="spotify.agg_playlist_history", )) 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-agg 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_agg, date=ds)) if not status: logger.error(f"Failed to append stats-agg for {d}") return None logger.info("Rebuilding stats-agg hyper table") status = rdb.execute(ql.load(QUERY_PATHS.make_agg_hyper)) if not status: logger.error("Failed to rebuild stats-agg hyper table") return None logger.success(f"Stats-agg updated to {date_list[-1]}") return date_list[-1].strftime("%Y-%m-%d")