""" OsservatorioReport — Report settimanale tabOsservatorio Genera Excel con nuovi prodotti in Osservatorio non ancora in GEMMA e invia via Outlook. Logica filtri: 1. codice > watermark (ultimo codice inviato) 2. anno >= 2025 3. IMDB non presente in GEMMA (gemma.parquet) Sorgenti: osservatorio.parquet + gemma.parquet Watermark salvato in osservatorio_tracker.db (tabella watermark, key='last_codice'). """ import argparse import json import sqlite3 from datetime import datetime from pathlib import Path import pandas as pd import openpyxl from openpyxl.styles import PatternFill, Font, Alignment from openpyxl.utils import get_column_letter from loguru import logger BASE_DIR = Path(__file__).parent # --------------------------------------------------------------------------- # Config # --------------------------------------------------------------------------- def load_config() -> dict: with open(BASE_DIR / "config.json", encoding="utf-8") as f: return json.load(f) # --------------------------------------------------------------------------- # Parquet L2 # --------------------------------------------------------------------------- def get_osservatorio_data(cfg: dict, codice_dal: int) -> pd.DataFrame: """ Legge osservatorio.parquet (codice > codice_dal, anno >= 2025). Esclude prodotti già presenti in gemma.parquet. """ base = cfg["parquet_base_path"] df = pd.read_parquet(f"{base}/osservatorio.parquet") df = df.rename(columns={"cast": "Cast"}) df["codice"] = pd.to_numeric(df["codice"], errors="coerce") df["anno"] = pd.to_numeric(df["anno"], errors="coerce") df = df[df["codice"] > codice_dal] df = df[df["anno"] >= 2025] df = df.sort_values(["anno", "titolo"]).reset_index(drop=True) logger.info(f"osservatorio.parquet: {len(df)} righe (codice > {codice_dal})") if df.empty: return df # IMDB già in GEMMA df_gemma = pd.read_parquet(f"{base}/gemma.parquet", columns=["A_COD_IMDB"]) gemma_imdbs = set(df_gemma["A_COD_IMDB"].dropna().unique()) logger.info(f"gemma.parquet: {len(gemma_imdbs)} IMDB presenti") n0 = len(df) df = df[~df["imdb"].isin(gemma_imdbs)] logger.debug(f"Dopo filtro GEMMA: {len(df)} (rimossi {n0 - len(df)})") return df.reset_index(drop=True) # --------------------------------------------------------------------------- # Watermark # --------------------------------------------------------------------------- def _tracker_conn(cfg: dict) -> sqlite3.Connection: path = BASE_DIR / cfg.get("osservatorio_tracker_db", "osservatorio_tracker.db") conn = sqlite3.connect(path) conn.execute(""" CREATE TABLE IF NOT EXISTS watermark ( key TEXT PRIMARY KEY, value TEXT NOT NULL ) """) conn.commit() return conn def get_last_codice(cfg: dict) -> int: """Legge il watermark (ultimo codice inviato). Fallback: config osservatorio_codice_inizio.""" conn = _tracker_conn(cfg) row = conn.execute("SELECT value FROM watermark WHERE key='last_codice'").fetchone() conn.close() if row: return int(row[0]) return int(cfg.get("osservatorio_codice_inizio", 75104)) def mark_watermark(cfg: dict, max_codice: int): conn = _tracker_conn(cfg) conn.execute( "INSERT OR REPLACE INTO watermark (key, value) VALUES ('last_codice', ?)", [str(max_codice)], ) conn.commit() conn.close() logger.info(f"Osservatorio watermark aggiornato: last_codice={max_codice}") # --------------------------------------------------------------------------- # Generazione Excel # --------------------------------------------------------------------------- HEADERS = ["stato", "tipologia", "titolo", "titolo_int", "anno", "paese", "genere", "produzione", "distribuzione", "regia", "Cast", "trama", "rete", "platform", "imdb", "codice"] COL_MAP = ["stato", "tipologia", "titolo", "titolo_int", "anno", "paese", "genere", "produzione", "distribuzione", "regia", "Cast", "trama", "rete", "platform", "imdb", "codice"] COL_WIDTHS = [12, 14, 36, 26, 7, 12, 18, 22, 22, 18, 18, 40, 10, 12, 14, 8] BLUE = PatternFill(start_color="DDEBF7", end_color="DDEBF7", fill_type="solid") FONT_NAME = "Aptos Narrow" def generate_excel(df: pd.DataFrame) -> Path: output_dir = BASE_DIR / "output" output_dir.mkdir(exist_ok=True) output_path = output_dir / f"OSSERVATORIO_{datetime.now().strftime('%Y%m%d')}.xlsx" wb = openpyxl.Workbook() ws = wb.active ws.title = "OSSERVATORIO" # Riga 1: intestazioni + autofilter for col, header in enumerate(HEADERS, 1): cell = ws.cell(row=1, column=col, value=header) cell.font = Font(name=FONT_NAME, bold=True, size=9) cell.fill = BLUE cell.alignment = Alignment(horizontal="center", vertical="center") ws.auto_filter.ref = f"A1:{get_column_letter(len(HEADERS))}1" ws.row_dimensions[1].height = 14 # Larghezze colonne for i, width in enumerate(COL_WIDTHS, 1): ws.column_dimensions[get_column_letter(i)].width = width # Dati (riga 2+) for row_idx, row in df.iterrows(): excel_row = row_idx + 2 for col_idx, col_name in enumerate(COL_MAP, 1): value = row[col_name] if pd.isna(value): value = None cell = ws.cell(row=excel_row, column=col_idx, value=value) cell.font = Font(name=FONT_NAME, size=9) wb.save(output_path) logger.info(f"Excel generato: {output_path.name} ({len(df)} righe)") return output_path # --------------------------------------------------------------------------- # Invio email # --------------------------------------------------------------------------- def send_email(excel_path, cfg: dict, n_rows: int): import win32com.client email_cfg = cfg["osservatorio_email"] subject = email_cfg["subject"] body = email_cfg.get("body", "") recipients = "; ".join(email_cfg["to"]) outlook = win32com.client.Dispatch("Outlook.Application") mail = outlook.CreateItem(0) mail.Subject = subject mail.Body = body mail.To = recipients if email_cfg.get("cc"): mail.CC = "; ".join(email_cfg["cc"]) if excel_path is not None: mail.Attachments.Add(str(excel_path.resolve())) mail.Send() logger.info(f"Email inviata — {'con allegato' if excel_path else 'senza allegato'} — TO: {len(email_cfg['to'])}") # --------------------------------------------------------------------------- # Entry point # --------------------------------------------------------------------------- def run(dry_run: bool = False): logs_dir = BASE_DIR / "logs" logs_dir.mkdir(exist_ok=True) logger.add( logs_dir / "osservatorio_report_{time:YYYY-MM-DD}.log", rotation="7 days", retention="30 days", level="DEBUG", ) mode = "[DRY-RUN] " if dry_run else "" logger.info(f"=== OsservatorioReport avviato {mode}===") cfg = load_config() codice_dal = get_last_codice(cfg) logger.info(f"Watermark: prodotti con codice > {codice_dal}") df = get_osservatorio_data(cfg, codice_dal) logger.info(f"Nuovi prodotti da inviare: {len(df)}") if df.empty: logger.info("Nessun nuovo prodotto questa settimana — invio email vuota") if not dry_run: send_email(None, cfg, 0) else: logger.info("[DRY-RUN] Email vuota NON inviata") logger.info(f"=== OsservatorioReport completato {mode}===") return excel_path = generate_excel(df) if dry_run: logger.info(f"[DRY-RUN] Email NON inviata — Excel generato in: {excel_path}") logger.info(f"[DRY-RUN] Destinatari TO: {cfg['osservatorio_email']['to']}") logger.info("[DRY-RUN] Watermark NON aggiornato") else: send_email(excel_path, cfg, len(df)) mark_watermark(cfg, int(df["codice"].max())) logger.info(f"=== OsservatorioReport completato {mode}===") if __name__ == "__main__": parser = argparse.ArgumentParser() parser.add_argument("--dry-run", action="store_true") args = parser.parse_args() run(dry_run=args.dry_run)