""" Parquet L3 aggregation — generazione emesso_agg.parquet + custom_agg.parquet. Chiamato da bootstrap.py (al sync) e da linker/launcher.py (avvio da admin panel). Prende app_root e server_root espliciti — nessuna dipendenza da costanti di bootstrap. """ import shutil from pathlib import Path _GENERALISTE = ('N1', 'N2', 'N3', 'C5', 'I1', 'MC', 'R4') _SQLITE_LAYER2_DIR = Path(r"I:\SOFTWARE\PYTHON_SRV\sqlite_layer2") _SERVER_ROOT_FALLBACK = Path(r"I:\SOFTWARE\PYTHON_SRV\PYTHON_LOCAL\MyICR_Suite") def l3_needs_rebuild(app_root: Path) -> bool: """True se i parquet L3 mancano o sono più vecchi di almeno un parquet L2 chiave.""" marker = app_root / "local_db" / "parquet_agg" / "emesso_agg.parquet" if not marker.exists(): return True l3_mtime = marker.stat().st_mtime local_parquet = app_root / "local_db" / "parquet" for name in ("emesso.parquet", "prodotti.parquet", "cast.parquet", "diritti.parquet", "imdb.parquet"): f = local_parquet / name if f.exists() and f.stat().st_mtime > l3_mtime: return True return False def sync_parquet_l3(app_root: Path, server_root: Path | None = None, progress_cb=None): """Garantisce che i parquet L3 aggregati siano aggiornati. Flusso: 1. Nessun rebuild necessario → exit subito. 2. Il server ha già i L3 aggiornati (mtime >= emesso.parquet locale) → pull. 3. Altrimenti → build locale → push su server → scrivi mirror_l2.sqlite. """ if server_root is None: server_root = _SERVER_ROOT_FALLBACK local_parquet = app_root / "local_db" / "parquet" agg_dir = app_root / "local_db" / "parquet_agg" server_agg = server_root / "local_db" / "parquet_agg" if not l3_needs_rebuild(app_root): return # Prova a pullare dal server se ha già i L3 aggiornati server_marker = server_agg / "emesso_agg.parquet" emesso_l2 = local_parquet / "emesso.parquet" if (server_marker.exists() and emesso_l2.exists() and server_marker.stat().st_mtime >= emesso_l2.stat().st_mtime): if progress_cb: progress_cb(83, "Aggiornamento dati aggregati dal server...") agg_dir.mkdir(parents=True, exist_ok=True) for name in ("emesso_agg.parquet", "custom_agg.parquet"): src = server_agg / name if src.exists(): shutil.copy2(src, agg_dir / name) return # Build locale if progress_cb: progress_cb(83, "Elaborazione dati in corso…") build_parquet_l3(local_parquet, agg_dir, progress_cb) # Push parquet_agg/ su server if server_agg.parent.exists(): if progress_cb: progress_cb(87, "Pubblicazione dati aggregati su server...") server_agg.mkdir(parents=True, exist_ok=True) for name in ("emesso_agg.parquet", "custom_agg.parquet"): src = agg_dir / name if src.exists(): shutil.copy2(src, server_agg / name) # Scrivi mirror_l2.sqlite (solo se la dir server è raggiungibile) if _SQLITE_LAYER2_DIR.exists(): if progress_cb: progress_cb(87, "Aggiornamento SQLite layer2...") try: write_mirror_l2_sqlite(agg_dir) except Exception as e: print(f"[parquet_l3] mirror_l2.sqlite fallito (non critico): {e}") def build_parquet_l3(local_parquet: Path, agg_dir: Path, progress_cb=None): """Legge i parquet L2 locali e scrive emesso_agg.parquet + custom_agg.parquet.""" import duckdb agg_dir.mkdir(parents=True, exist_ok=True) em_f = str(local_parquet / "emesso.parquet").replace('\\', '/') prod_f = str(local_parquet / "prodotti.parquet").replace('\\', '/') cast_f = str(local_parquet / "cast.parquet").replace('\\', '/') imdb_f = str(local_parquet / "imdb.parquet").replace('\\', '/') dir_f = str(local_parquet / "diritti.parquet").replace('\\', '/') con = duckdb.connect() nets = ', '.join(f"'{n}'" for n in _GENERALISTE) if progress_cb: progress_cb(84, "Elaborazione dati in corso… (emesso)") emesso_out = str(agg_dir / "emesso_agg.parquet").replace('\\', '/') con.execute(f""" COPY ( WITH base AS ( SELECT CAST(prodotto AS VARCHAR) AS RIFER, rete, data_emissione, ora_inizio, fascia, audience, share, CASE WHEN rete IN ({nets}) THEN 'G' ELSE 'T' END AS tipo_rete FROM read_parquet('{em_f}') ) SELECT RIFER, tipo_rete, FIRST(rete ORDER BY data_emissione ASC, ora_inizio ASC) AS rete_first, FIRST(data_emissione ORDER BY data_emissione ASC, ora_inizio ASC) AS data_first, FIRST(ora_inizio ORDER BY data_emissione ASC, ora_inizio ASC) AS ora_first, FIRST(fascia ORDER BY data_emissione ASC, ora_inizio ASC) AS fascia_first, FIRST(audience ORDER BY data_emissione ASC, ora_inizio ASC) AS aud_first, FIRST(share ORDER BY data_emissione ASC, ora_inizio ASC) AS sha_first, LAST(rete ORDER BY data_emissione ASC, ora_inizio ASC) AS rete_last, LAST(data_emissione ORDER BY data_emissione ASC, ora_inizio ASC) AS data_last, LAST(ora_inizio ORDER BY data_emissione ASC, ora_inizio ASC) AS ora_last, LAST(fascia ORDER BY data_emissione ASC, ora_inizio ASC) AS fascia_last, LAST(audience ORDER BY data_emissione ASC, ora_inizio ASC) AS aud_last, LAST(share ORDER BY data_emissione ASC, ora_inizio ASC) AS sha_last FROM base GROUP BY RIFER, tipo_rete ) TO '{emesso_out}' (FORMAT PARQUET) """) if progress_cb: progress_cb(86, "Elaborazione dati in corso… (anagrafica)") custom_out = str(agg_dir / "custom_agg.parquet").replace('\\', '/') con.execute(f""" COPY ( WITH prod AS ( SELECT * FROM read_parquet('{prod_f}') QUALIFY ROW_NUMBER() OVER (PARTITION BY prodotto ORDER BY edizione) = 1 ), imdb AS ( SELECT codice, MIN(riferimento_imdb) AS riferimento_imdb FROM read_parquet('{imdb_f}') GROUP BY codice ), reg AS ( SELECT prodotto, cognome AS REGISTA_COGN, nome AS REGISTA_NOME FROM read_parquet('{cast_f}') WHERE ruolo = 'FRE' QUALIFY ROW_NUMBER() OVER (PARTITION BY prodotto ORDER BY progr_cast) = 1 ), a1 AS ( SELECT prodotto, TRIM(COALESCE(nome,'') || ' ' || COALESCE(cognome,'')) AS ATTORE1 FROM read_parquet('{cast_f}') WHERE ruolo='C001' AND progr_cast=1 ), a2 AS ( SELECT prodotto, TRIM(COALESCE(nome,'') || ' ' || COALESCE(cognome,'')) AS ATTORE2 FROM read_parquet('{cast_f}') WHERE ruolo='C001' AND progr_cast=2 ), a3 AS ( SELECT prodotto, TRIM(COALESCE(nome,'') || ' ' || COALESCE(cognome,'')) AS ATTORE3 FROM read_parquet('{cast_f}') WHERE ruolo='C001' AND progr_cast=3 ), distr AS ( SELECT prod AS prodotto, MIN(ragsoc_distr) AS DF_RAGSOC_DISTR FROM read_parquet('{dir_f}') GROUP BY prod ) SELECT CAST(p.prodotto AS VARCHAR) AS RIFER, p.tipologia AS TIPOL, p.titolo_originale AS "TO", p.titolo_italiano AS TI, p.anno_produzione AS ANNO, p.paesi_produzione1 AS PAESE, p.genere1 AS GENERE1, p.num_episodi AS EPIS, p.durata AS DUR, p.superserie_descr AS SUPERSERIE, CASE p.veg WHEN 'EVER GREEN' THEN 'EV.GR' ELSE p.veg END AS VEG_FREE, i.riferimento_imdb AS IMDB_CODICE, r.REGISTA_COGN, r.REGISTA_NOME, a1.ATTORE1, a2.ATTORE2, a3.ATTORE3, d.DF_RAGSOC_DISTR FROM prod p LEFT JOIN imdb i ON i.codice = p.prodotto LEFT JOIN reg r ON r.prodotto = p.prodotto LEFT JOIN a1 ON a1.prodotto = p.prodotto LEFT JOIN a2 ON a2.prodotto = p.prodotto LEFT JOIN a3 ON a3.prodotto = p.prodotto LEFT JOIN distr d ON d.prodotto = p.prodotto ) TO '{custom_out}' (FORMAT PARQUET) """) con.close() def write_mirror_l2_sqlite(agg_dir: Path): """Scrive emesso_agg e custom_agg in mirror_l2.sqlite sul server (uso personale Mauro).""" import duckdb import sqlite3 as _sl mirror_path = _SQLITE_LAYER2_DIR / "mirror_l2.sqlite" con = duckdb.connect() emesso_f = str(agg_dir / "emesso_agg.parquet").replace('\\', '/') custom_f = str(agg_dir / "custom_agg.parquet").replace('\\', '/') df_emesso = con.execute(f"SELECT * FROM read_parquet('{emesso_f}')").df() df_custom = con.execute(f"SELECT * FROM read_parquet('{custom_f}')").df() con.close() sl_con = _sl.connect(str(mirror_path)) df_emesso.to_sql("emesso_agg", sl_con, if_exists='replace', index=False) df_custom.to_sql("custom_agg", sl_con, if_exists='replace', index=False) sl_con.execute("PRAGMA journal_mode=WAL") sl_con.close()