from __future__ import annotations from datetime import date, timedelta import pandas as pd from loguru import logger from ..config import QUERY_PATHS, TABLES, WAREHOUSES from ..db import get_snw from ..db.queries import QueryLoader from .api import fetch_rates def _write_to_snowflake(df: pd.DataFrame) -> int: """Write a rates DataFrame to Snowflake via staging + merge. Returns the number of rows merged. """ if df.empty: logger.info("No rates to write") return 0 snw = get_snw(WAREHOUSES["small"]) ql = QueryLoader() # Ensure main table exists (safe for repeated runs) snw.execute(ql.load(QUERY_PATHS.snw_ensure)) # Create a temp staging table, load data, then merge into main staging_table = f"{TABLES['snw']}_staging" snw.execute(ql.load(QUERY_PATHS.snw_create_staging, staging=staging_table)) snw.write(df, staging_table) logger.debug(f"Wrote {len(df)} rows to staging table {staging_table}") merge_sql = ql.load(QUERY_PATHS.snw_upsert, staging=staging_table) snw.execute(merge_sql) logger.info(f"Merged {len(df)} rate rows into {TABLES['snw']}") return len(df) def run_backfill(start_date: date | None = None) -> None: """Fetch historical rates and load into Snowflake.""" if start_date is None: start_date = date.today() - timedelta(days=5 * 365) end_date = date.today() - timedelta(days=1) logger.info(f"Starting backfill from {start_date} to {end_date}") df = fetch_rates(start_date, end_date) merged = _write_to_snowflake(df) logger.success(f"Backfill complete — {merged} rate rows loaded") def run_daily() -> None: """Fetch the most recent rates and upsert into Snowflake.""" snw = get_snw(WAREHOUSES["small"]) ql = QueryLoader() # Find the latest date already in Snowflake result = snw.query(ql.load(QUERY_PATHS.snw_get_cutoff)) max_date = result["max_date"][0] if max_date is None: logger.warning("Table is empty — run backfill first") return if isinstance(max_date, str): max_date = date.fromisoformat(max_date) start_date = max_date + timedelta(days=1) end_date = date.today() - timedelta(days=1) if start_date > end_date: logger.info("Already up to date — no new rates to fetch") return logger.info(f"Fetching daily rates from {start_date} to {end_date}") df = fetch_rates(start_date, end_date) merged = _write_to_snowflake(df) logger.success(f"Daily ingest complete — {merged} rate rows loaded")