""" DataHub v2 — CLI entry point Uso: python main.py weekly # tutti i Parquet weekly python main.py daily # gemma + osservatorio python main.py weekly --dry-run # logga conteggi, non scrive python main.py prodotti # singolo domain (debug) python main.py diritti # singolo domain (debug) """ import sys import os # Python embedded non include la dir corrente nel path — va aggiunto esplicitamente sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) import argparse import duckdb from loguru import logger from config.settings import ensure_dirs, get_log_path, LOG_FILE # Mappa singolo domain → modulo DOMAIN_MAP = { 'reti_cluster': 'domains.reti_cluster', 'prodotti': 'domains.prodotti', 'cast': 'domains.cast', 'imdb': 'domains.imdb', 'veg': 'domains.veg', 'scelte_rete': 'domains.scelte_rete', 'boxoffice': 'domains.boxoffice', 'boxoffice_dettaglio': 'domains.boxoffice_dettaglio', 'osservatorio': 'domains.osservatorio', 'diritti': 'domains.diritti', 'diritti_enrich': 'domains.diritti_enrich', 'emesso': 'domains.emesso', 'emesso_enrich': 'domains.emesso_enrich', 'gemma': 'domains.gemma', } def _setup_logging(prefix: str) -> None: log_path = LOG_FILE if LOG_FILE else get_log_path(prefix) logger.remove() logger.add(sys.stderr, level='INFO', format='{time:HH:mm:ss} | {level:<8} | {message}') logger.add(log_path, level='DEBUG', mode='a', format='{time:YYYY-MM-DD HH:mm:ss} | {level:<8} | {message}', encoding='utf-8') def _run_single(domain: str, dry_run: bool) -> None: import importlib mod_path = DOMAIN_MAP.get(domain) if not mod_path: logger.error(f"Domain '{domain}' non riconosciuto. Disponibili: {list(DOMAIN_MAP)}") sys.exit(1) mod = importlib.import_module(mod_path) with duckdb.connect() as con: logger.info(f"Avvio domain: {domain}") result = mod.run(con) logger.success(f"{domain}: {result:,} righe") def main() -> None: parser = argparse.ArgumentParser( description='DataHub v2 — pipeline CSV → Parquet' ) parser.add_argument( 'target', help=( 'Cosa eseguire: "weekly", "daily", ' 'oppure nome di un singolo domain (es. "prodotti")' ) ) parser.add_argument( '--dry-run', action='store_true', help='Esegue senza scrivere file (solo log conteggi)' ) args = parser.parse_args() ensure_dirs() _setup_logging(args.target) logger.info(f"DataHub v2 — target={args.target} dry_run={args.dry_run}") if args.target == 'weekly': from pipelines import weekly results = weekly.run(dry_run=args.dry_run) _print_summary(results) elif args.target == 'daily': from pipelines import daily results = daily.run(dry_run=args.dry_run) _print_summary(results) elif args.target in DOMAIN_MAP: _run_single(args.target, dry_run=args.dry_run) else: logger.error( f"Target '{args.target}' non valido. " f"Usa: weekly | daily | {' | '.join(DOMAIN_MAP)}" ) sys.exit(1) def _print_summary(results: dict) -> None: logger.info("─" * 50) logger.info("RIEPILOGO") logger.info("─" * 50) ok = all(v >= 0 for v in results.values()) for name, count in results.items(): if count < 0: logger.error(f" {name:<25} ERRORE") else: logger.success(f" {name:<25} {count:>10,} righe") logger.info("─" * 50) if ok: logger.success("Pipeline completata con successo.") else: logger.warning("Pipeline completata con errori (vedi log sopra).") if __name__ == '__main__': main()