""" Pipeline daily: produce gemma.parquet e osservatorio.parquet Post-step: esporta gemma.db (SQLite) per MediaTrack che usa ATTACH DATABASE. """ import os import duckdb from loguru import logger from config.settings import out_parquet, OUTPUT_DIR from domains import gemma, osservatorio DAILY_DOMAINS = [ ('gemma', gemma), ('osservatorio', osservatorio), ] def _export_gemma_sqlite() -> None: """Esporta gemma.parquet → gemma.db (SQLite) via DuckDB SQLite extension.""" parquet = out_parquet('gemma').replace('\\', '/') db_path = os.path.join(OUTPUT_DIR, 'gemma.db').replace('\\', '/') with duckdb.connect() as con: con.execute("INSTALL sqlite; LOAD sqlite") con.execute(f"ATTACH '{db_path}' AS _gemma_db (TYPE SQLITE)") con.execute("CREATE OR REPLACE TABLE _gemma_db.gemma AS SELECT * FROM read_parquet('" + parquet + "')") con.execute("DETACH _gemma_db") logger.success(f"[daily] gemma.db esportato → {db_path}") def run(dry_run: bool = False) -> dict: """ Esegue la pipeline daily. Restituisce un dict {nome: count}. """ results = {} with duckdb.connect() as con: for name, mod in DAILY_DOMAINS: logger.info(f"[daily] avvio domain: {name}") try: result = mod.run(con) results[name] = result logger.success(f"[daily] {name}: {result:,} righe") except Exception as exc: logger.error(f"[daily] {name} FALLITO: {exc}") results[name] = -1 if results.get('gemma', -1) >= 0 and not dry_run: try: _export_gemma_sqlite() except Exception as exc: logger.error(f"[daily] gemma.db export FALLITO: {exc}") return results