import os import re from pathlib import Path from typing import Literal import pandas as pd from loguru import logger from ..config.paths import PATHS, HEADLESS_CHROME_PATH from ..config.countries import COUNTRY_DICT, LABEL_DICT from ..config.email import MULTI_MARKET_REPORTS from ..email.renderer import EmailRenderer, ReportType from ..email.screenshot import ScreenshotService from ..email.composer import EmailComposer from ..email.sender import send_batch, chunk_recipients from .distro import DistroService Mode = Literal["dev", "prod"] class EmailService: def __init__(self): self.renderer = EmailRenderer() self.composer = EmailComposer() self.distro_service = DistroService() self.screenshot: ScreenshotService | None = None self.spotify_cid: str = "" self._cover_cache: dict[str, str] = {} def send_all(self, mode: Mode = "prod", affiliate: str | None = None, to: str | None = None, cleanup: bool = True) -> None: data_all = self._load_latest_csv() if data_all.empty: logger.warning("No data to send") return self._init_screenshot_service() to_override = [addr.strip() for addr in to.split(",") if addr.strip()] if to else None if affiliate: self._send_single_affiliate(data_all, mode, affiliate, to_override) else: self._cover_cache = self.screenshot.download_covers_batch( data_all["artwork_url"].dropna().tolist() ) self._send_country_reports(data_all, mode, to_override) self._send_label_reports(data_all, mode, to_override) self._send_cea_report(data_all, mode, to_override) self._send_multi_market_reports(data_all, mode, to_override) self.screenshot.cleanup_covers() if cleanup: self._cleanup_daily_files() def _send_single_affiliate(self, data_all: pd.DataFrame, mode: Mode, affiliate: str, to_override: list[str] | None = None) -> None: distro = self.distro_service.get_distro(mode) if affiliate == "CEA": self._send_cea_report(data_all, mode, to_override) return if affiliate in LABEL_DICT: info = LABEL_DICT[affiliate] logger.info(f"Processing {info['name']}") rep_owners = info["rep_owners"] data = data_all.query("rep_owner in @rep_owners").reset_index() if not data.empty: self._send_report(data, affiliate, to_override or distro.get(affiliate, []), "lb", mode) return if affiliate in COUNTRY_DICT: info = COUNTRY_DICT[affiliate] logger.info(f"Processing {info['name']}") if info.get("generate_cc_mail"): data_cc = data_all.query("country_code == @affiliate").reset_index() if not data_cc.empty: self._send_report(data_cc, affiliate, to_override or distro.get(affiliate, []), "cc", mode) if info.get("generate_ro_mail"): affiliates = info.get("affiliate_territories", [affiliate]) data_ro = data_all.query( "owner_abbreviation == @affiliate and country_code not in @affiliates and country_code != 'WW'" ).reset_index() if not data_ro.empty: data_ro = data_ro.sort_values(by=["isrc_cd", "country_code"]).reset_index(drop=True) self._send_report(data_ro, affiliate, to_override or distro.get(f"x{affiliate}", []), "ro", mode) return logger.error(f"Unknown affiliate: {affiliate}") def _init_screenshot_service(self) -> None: self.screenshot = ScreenshotService( output_path=PATHS.html_tmp, browser_executable=HEADLESS_CHROME_PATH, ) self.spotify_cid = self.screenshot.prepare_spotify_logo() def _send_country_reports(self, data_all: pd.DataFrame, mode: Mode, to_override: list[str] | None = None) -> None: distro = self.distro_service.get_distro(mode) for country, info in COUNTRY_DICT.items(): if not info.get("generate_cc_mail"): continue logger.info(f"Processing {info['name']}") data_cc = data_all.query("country_code == @country").reset_index() if not data_cc.empty: self._send_report( data=data_cc, affiliate=country, recipients=to_override or distro.get(country, []), report_type="cc", mode=mode, ) if not info.get("generate_ro_mail"): continue affiliates = info.get("affiliate_territories", [country]) data_ro = data_all.query( "owner_abbreviation == @country and country_code not in @affiliates and country_code != 'WW'" ).reset_index() if not data_ro.empty: data_ro = data_ro.sort_values(by=["isrc_cd", "country_code"]).reset_index(drop=True) self._send_report( data=data_ro, affiliate=country, recipients=to_override or distro.get(f"x{country}", []), report_type="ro", mode=mode, ) def _send_label_reports(self, data_all: pd.DataFrame, mode: Mode, to_override: list[str] | None = None) -> None: distro = self.distro_service.get_distro(mode) for label, info in LABEL_DICT.items(): logger.info(f"Processing {info['name']}") rep_owners = info["rep_owners"] data_label = data_all.query("rep_owner in @rep_owners").reset_index() if not data_label.empty: self._send_report( data=data_label, affiliate=label, recipients=to_override or distro.get(label, []), report_type="lb", mode=mode, ) @staticmethod def _filter_cea_data(data_all: pd.DataFrame) -> pd.DataFrame: """Filter data for CEA aggregated reports.""" internal = [c for c, i in COUNTRY_DICT.items() if i.get("is_cea")] external = [c for c, i in COUNTRY_DICT.items() if not i.get("is_cea")] data = data_all.query("country_code != 'WW'") return data.query( "(owner_abbreviation in @internal and n_markets >= 5) or " "(owner_abbreviation in @external and country_code in @internal and n_markets_cea >= 8)" ).reset_index(drop=True) def _send_cea_report(self, data_all: pd.DataFrame, mode: Mode, to_override: list[str] | None = None) -> None: distro = self.distro_service.get_distro(mode) data = self._filter_cea_data(data_all) if data.empty: return logger.info("Processing CEA aggregated report") data = data.sort_values(by=["owner_abbreviation", "isrc_cd", "country_code"]).reset_index(drop=True) self._send_report( data=data, affiliate="CEA", recipients=to_override or distro.get("CEA", []), report_type="ro", mode=mode, show_owner=True, ) def _send_multi_market_reports(self, data_all: pd.DataFrame, mode: Mode, to_override: list[str] | None = None) -> None: """Send all configured multi-market reports.""" distro = self.distro_service.get_distro(mode) for report in MULTI_MARKET_REPORTS: markets = report["markets"] distro_key = report["distro_key"] recipients = to_override or distro.get(distro_key, []) self._send_multi_market(data_all, markets, recipients, mode) def _send_multi_market( self, data_all: pd.DataFrame, markets: list[str], recipients: list[str], mode: Mode, ) -> None: """Build and send a multi-market report for the given markets.""" data = data_all.query("country_code in @markets").reset_index(drop=True) # Only keep tracks that appear in at least 2 of the specified markets isrc_market_counts = data.groupby("isrc_cd")["country_code"].nunique() qualifying_isrcs = isrc_market_counts[isrc_market_counts >= 2].index data = data[data["isrc_cd"].isin(qualifying_isrcs)].reset_index(drop=True) if data.empty: logger.warning(f"No data for markets: {markets}") return market_names = ", ".join( COUNTRY_DICT.get(m, {}).get("name", m) for m in markets ) subject = f"WHATS COOKIN, {market_names} edition" logger.info(f"Building multi-market email: {subject} ({len(data)} track/country rows)") body_html, cids, attachments = self._build_email_content(data, "mm") if not recipients: return emails = [ self.composer.compose( subject=subject, recipients=chunk, body_html=body_html, image_attachments=attachments, mode=mode, ) for chunk in chunk_recipients(recipients) ] send_batch(emails) self.screenshot.cleanup() logger.success(f"Sent {subject}") def send_multi_market( self, markets: list[str], mode: Mode = "prod", to: str | None = None, ) -> None: """Send a multi-market report for an arbitrary combination of markets (CLI entry point).""" data_all = self._load_latest_csv() if data_all.empty: logger.warning("No data to send") return self._init_screenshot_service() to_override = [addr.strip() for addr in to.split(",") if addr.strip()] if to else None recipients = to_override or [] if not recipients: distro = self.distro_service.get_distro(mode) recipients = distro.get("CEA", []) self._send_multi_market(data_all, markets, recipients, mode) self.screenshot.cleanup_covers() def _send_report( self, data: pd.DataFrame, affiliate: str, recipients: list[str], report_type: ReportType, mode: Mode, show_owner: bool = False, ) -> None: if data.empty or not recipients: return subject = self._get_subject(affiliate, report_type) logger.info(f"Building email: {subject} ({len(data)} tracks)") body_html, cids, attachments = self._build_email_content(data, report_type, show_owner=show_owner) emails = [ self.composer.compose( subject=subject, recipients=chunk, body_html=body_html, image_attachments=attachments, mode=mode, ) for chunk in chunk_recipients(recipients) ] send_batch(emails) self.screenshot.cleanup() logger.success(f"Sent {subject}") def _build_email_content( self, data: pd.DataFrame, report_type: ReportType, show_owner: bool = False, ) -> tuple[str, list[str], list[tuple[Path, str]]]: """Build email content for any report type (cc, ro, lb, mm).""" multi_market = report_type in ("ro", "lb", "mm") if multi_market: # Split: ISRCs with 1 country use single-market cards, 2+ use multi-market cards isrc_counts = data.groupby("isrc_cd")["country_code"].nunique() single_isrcs = isrc_counts[isrc_counts == 1].index multi_isrcs = isrc_counts[isrc_counts >= 2].index data_single = data[data["isrc_cd"].isin(single_isrcs)].reset_index(drop=True) data_multi = data[data["isrc_cd"].isin(multi_isrcs)].reset_index(drop=True) parsed_single = self.renderer.parse_report(data_single, report_type) if not data_single.empty else None parsed_multi = self.renderer.parse_multi_market_report(data_multi, show_owner=show_owner, report_type=report_type) if not data_multi.empty else None # Group tracks by rep_owner when show_owner is enabled (CEA aggregated report) if show_owner: if parsed_multi is not None and not parsed_multi.empty: parsed_multi = parsed_multi.sort_values( by=["owner_abbreviation", "total_streams"], ascending=[True, False], ).reset_index(drop=True) if parsed_single is not None and not parsed_single.empty: parsed_single = parsed_single.sort_values( by=["owner_abbreviation", "sum_streams_spotify"], ascending=[True, False], ).reset_index(drop=True) else: parsed_single = self.renderer.parse_report(data, report_type) parsed_multi = None # Ensure covers are downloaded for all tracks all_urls = [] if parsed_single is not None: all_urls += parsed_single["artwork_url"].dropna().tolist() if parsed_multi is not None: all_urls += parsed_multi["artwork_url"].dropna().tolist() missing = [u for u in all_urls if u and str(u) != "nan" and u not in self._cover_cache] if missing: new_covers = self.screenshot.download_covers_batch(missing) self._cover_cache.update(new_covers) all_sections = "" cids: list[str] = [] attachments: list[tuple[Path, str]] = [] track_index = 0 # Render multi-market cards first (higher priority / more markets) if parsed_multi is not None: for _, row in parsed_multi.iterrows(): cover_path = self.screenshot.get_cover_path(row["artwork_url"], self._cover_cache) row_html = self.renderer.render_multi_market_section(row, cover_path) n_country_rows = row["n_country_rows"] image_cid = self.screenshot.screenshot_multi_market_row(track_index, row_html, n_country_rows) card_content_height = 80 + 25 + (n_country_rows * 25) + 10 total_height = 4 + card_content_height + 16 img_height = total_height * 0.75 section_html = self.composer.build_section_html( insights_link=row["insights_link"], spotify_link=row["spotify_link"], image_cid=image_cid, spotify_cid=self.spotify_cid, img_height=img_height, ) all_sections += section_html cids.append(image_cid) attachments.append((self.screenshot.get_track_image_path(track_index), image_cid)) track_index += 1 # Render single-market cards if parsed_single is not None: for _, row in parsed_single.iterrows(): cover_path = self.screenshot.get_cover_path(row["artwork_url"], self._cover_cache) row_html = self.renderer.render_section(row, cover_path) image_cid = self.screenshot.screenshot_row(track_index, row_html) section_html = self.composer.build_section_html( insights_link=row["insights_link"], spotify_link=row["spotify_link"], image_cid=image_cid, spotify_cid=self.spotify_cid, ) all_sections += section_html cids.append(image_cid) attachments.append((self.screenshot.get_track_image_path(track_index), image_cid)) track_index += 1 disclaimers = { "spotify": self._make_disclaimer(data, "report_date", "spotify_chart_date"), "tiktok": self._make_disclaimer(data, "report_date", "tiktok_rank_date"), "meta": self._make_disclaimer(data, "report_date", "meta_rank_date"), "shazam": self._make_disclaimer(data, "report_date", "shazam_chart_date"), } # Use multi-market legend if any multi-market cards were rendered legend_cid = self.screenshot.screenshot_legend(disclaimers, multi_market=(parsed_multi is not None)) report_date = data["report_date"].max().strftime("%Y-%m-%d") body_html = self.composer.build_body_html(all_sections, legend_cid, report_date=report_date) attachments.append((self.screenshot.get_legend_path(), legend_cid)) attachments.append((self.screenshot.get_spotify_logo_path(), self.spotify_cid)) return body_html, cids, attachments @staticmethod def _get_subject(affiliate: str, report_type: ReportType) -> str: if report_type == "cc": if affiliate == "WW": return "WHATS COOKIN around the world 🌎 🍳" info = COUNTRY_DICT.get(affiliate, {}) return f"WHATS COOKIN in {info.get('flag', '')} {info.get('emoji', '')}" if report_type == "lb": label_info = LABEL_DICT.get(affiliate, {}) return f"WHATS COOKIN 🍳, {label_info.get('name', affiliate)} edition" if affiliate == "CEA": return "WHATS COOKIN, CEA aggregated edition" country_info = COUNTRY_DICT.get(affiliate, {}) return f"WHATS COOKIN around the world, SME {country_info.get('name', affiliate)} edition" @staticmethod def _make_disclaimer(data: pd.DataFrame, wc_col: str, dsp_col: str) -> str: wc_d = data[wc_col].min() dsp_d = data[dsp_col].min() if pd.isna(dsp_d) or wc_d <= dsp_d: return "" day = dsp_d.day suffix = {1: "st", 2: "nd", 3: "rd"}.get(day % 10 if day % 100 not in [11, 12, 13] else 0, "th") return f"(*data for {dsp_d.strftime('%b')} {day}{suffix})" def _load_latest_csv(self) -> pd.DataFrame: date_cols = [ "report_date", "track_release_date", "tiktok_rank_date", "meta_rank_date", "spotify_chart_date", "shazam_chart_date", ] files = [f for f in os.listdir(PATHS.reports) if re.match(r"\d{4}-\d{2}-\d{2}\.csv", f)] if not files: return pd.DataFrame() latest = max(files) filepath = PATHS.reports / latest data = pd.read_csv(filepath, parse_dates=date_cols) return data.drop_duplicates() def _cleanup_daily_files(self) -> None: for f in os.listdir(PATHS.reports): if re.match(r"\d{4}-\d{2}-\d{2}\.csv", f): (PATHS.reports / f).unlink()