""" Query specifiche Linker + operazioni su LINKER_REPOSITORY """ import logging import sqlite3 import pandas as pd from datetime import datetime from pathlib import Path from typing import Optional, List, Dict from app.modules.linker.src.database.parquet_db import ParquetCatalog logger = logging.getLogger(__name__) _UPSERT_SQL = """ INSERT INTO LINKER_REPOSITORY (FORNITORE, TITOLO, TIPO_TITOLO, ANNO, REGISTA, RIFER, COD_IMDB, data_creazione, data_update) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(FORNITORE, TITOLO) DO UPDATE SET RIFER = excluded.RIFER, COD_IMDB = excluded.COD_IMDB, data_update = excluded.data_update """ class LinkerDatabase: """ Wrapper query specifiche per il modulo Linker. - CUSTOM_ANAGR_ORIZ / IMDB : ParquetCatalog - LINKER_REPOSITORY : SQLite (sqlite_con) - OTT write methods : server_con primary, sqlite_con mirror """ def __init__(self, catalog: ParquetCatalog, sqlite_con: sqlite3.Connection, server_con: sqlite3.Connection | None = None): self.catalog = catalog self.sqlite_con = sqlite_con self.server_con = server_con # primary write target for OTT tables self.last_server_write_ok: bool = True # consultato dalla GUI dopo ogni save OTT # ------------------------------------------------------------------ # # Dual-write helpers (server primary, local mirror best-effort) # # ------------------------------------------------------------------ # _SERVER_WRITE_TIMEOUT = 15.0 # secondi — oltre questo scrive solo locale e avvisa def _server_write(self, fn) -> bool: """Esegue fn(server_con) con timeout. Aggiorna last_server_write_ok. Ritorna True se OK.""" import threading done = threading.Event() error = [None] def _run(): try: fn(self.server_con) except Exception as e: error[0] = e finally: done.set() threading.Thread(target=_run, daemon=True).start() if not done.wait(timeout=self._SERVER_WRITE_TIMEOUT): logger.warning("[OTT server write] timeout dopo %.0fs — dato salvato solo in locale", self._SERVER_WRITE_TIMEOUT) self.last_server_write_ok = False return False if error[0]: logger.error(f"[OTT server write] {error[0]}") self.last_server_write_ok = False return False self.last_server_write_ok = True return True def _db_execute(self, sql: str, params: tuple = ()) -> bool: """Server prima (con timeout), poi mirror locale. Ritorna True se server OK.""" srv_ok = True if self.server_con is not None: srv_ok = self._server_write(lambda c: (c.execute(sql, params), c.commit())) self.sqlite_con.execute(sql, params) self.sqlite_con.commit() return srv_ok def _db_executemany(self, sql: str, params_list: list) -> bool: """executemany variant di _db_execute. Ritorna True se server OK.""" if not params_list: return True srv_ok = True if self.server_con is not None: pl = list(params_list) srv_ok = self._server_write(lambda c: (c.executemany(sql, pl), c.commit())) self.sqlite_con.executemany(sql, params_list) self.sqlite_con.commit() return srv_ok # ------------------------------------------------------------------ # # LINKER_REPOSITORY — lettura # # ------------------------------------------------------------------ # def get_fornitori(self) -> pd.DataFrame: try: return pd.read_sql_query( "SELECT DISTINCT FORNITORE FROM LINKER_REPOSITORY " "WHERE FORNITORE IS NOT NULL ORDER BY FORNITORE", self.sqlite_con, ) except Exception as e: logger.error(f"get_fornitori: {e}") return pd.DataFrame(columns=['FORNITORE']) def guess_fornitore(self, titles: list, min_matches: int = 2) -> tuple: """ Cerca di identificare il fornitore più probabile confrontando i titoli forniti con LINKER_REPOSITORY. Restituisce (fornitore, n_matches) o (None, 0) se non trovato. Usa LOWER() su entrambi i lati per confronto case-insensitive. """ sample = [t.strip().lower() for t in titles if t.strip()][:10] if not sample: return None, 0 placeholders = ','.join(['?' for _ in sample]) try: df = pd.read_sql_query( f"""SELECT FORNITORE, COUNT(*) AS cnt FROM LINKER_REPOSITORY WHERE LOWER(TITOLO) IN ({placeholders}) AND FORNITORE IS NOT NULL GROUP BY FORNITORE ORDER BY cnt DESC LIMIT 1""", self.sqlite_con, params=sample, ) if df.empty: return None, 0 row = df.iloc[0] cnt = int(row['cnt']) if cnt < min_matches: return None, 0 return str(row['FORNITORE']), cnt except Exception as e: logger.error(f"guess_fornitore: {e}") return None, 0 def get_repository_by_fornitore(self, fornitore: str) -> pd.DataFrame: try: return pd.read_sql_query( "SELECT * FROM LINKER_REPOSITORY WHERE UPPER(FORNITORE) = UPPER(?)", self.sqlite_con, params=(fornitore,), ) except Exception as e: logger.error(f"get_repository_by_fornitore: {e}") return pd.DataFrame() # ------------------------------------------------------------------ # # LINKER_REPOSITORY — scrittura # # ------------------------------------------------------------------ # def publish_to_server(self, server_path: str) -> tuple[int, int, int]: """ Upsert su server dei record LINKER_REPOSITORY locali che differiscono dal server. Confronta FORNITORE+TITOLO+RIFER+COD_IMDB — inserisce nuovi e aggiorna cambiati. Ritorna (upsertati, totale_locale, invariati). """ import sqlite3 as _sq3 local_df = pd.read_sql_query("SELECT * FROM LINKER_REPOSITORY", self.sqlite_con) if local_df.empty: return 0, 0, 0 srv_con = _sq3.connect(server_path) try: def _s(v) -> str: return str(v).strip() if v is not None and str(v) not in ('', 'nan', 'None') else '' try: srv_df = pd.read_sql_query( "SELECT FORNITORE, TITOLO, RIFER, COD_IMDB FROM LINKER_REPOSITORY", srv_con ) srv_index = { (_s(r['FORNITORE']).upper(), _s(r['TITOLO']).upper()): r for _, r in srv_df.iterrows() } except Exception: srv_index = {} to_upsert = [] for _, row in local_df.iterrows(): key = (_s(row.get('FORNITORE')).upper(), _s(row.get('TITOLO')).upper()) srv_row = srv_index.get(key) if srv_row is None: to_upsert.append(row) else: if (_s(row.get('RIFER')) != _s(srv_row.get('RIFER')) or _s(row.get('COD_IMDB')) != _s(srv_row.get('COD_IMDB'))): to_upsert.append(row) upserted = 0 cur = srv_con.cursor() for row in to_upsert: try: cur.execute(_UPSERT_SQL, ( row.get('FORNITORE'), row.get('TITOLO'), row.get('TIPO_TITOLO'), row.get('ANNO'), row.get('REGISTA'), row.get('RIFER'), row.get('COD_IMDB'), row.get('data_creazione'), row.get('data_update'), )) upserted += 1 except Exception as e: logger.error(f"publish_to_server upsert '{row.get('TITOLO')}': {e}") srv_con.commit() invariati = len(local_df) - upserted logger.info(f"publish_to_server: {upserted} upsertati, {invariati} invariati su {len(local_df)} locali") return upserted, len(local_df), invariati finally: srv_con.close() def save_to_repository(self, records: List[Dict], fornitore: Optional[str]) -> tuple[int, int]: """Salva in LINKER_REPOSITORY solo i record nuovi o cambiati. Ritorna (nuovi, aggiornati).""" if not records: return 0, 0 # Carica stato attuale per il fornitore — fonte di verità per il diff df_repo = self.get_repository_by_fornitore(fornitore or '') repo_index = { str(r.get('TITOLO') or '').upper(): { 'RIFER': str(r.get('RIFER') or '').strip(), 'COD_IMDB': str(r.get('COD_IMDB') or '').strip(), } for _, r in df_repo.iterrows() } now = datetime.now().isoformat(sep=' ') nuovi = aggiornati = 0 params_list = [] for rec in records: titolo = str(rec.get('TITOLO') or '').strip() new_rifer = str(rec.get('RIFER') or '').strip() new_imdb = str(rec.get('IMDB_CODICE') or '').strip() key = titolo.upper() existing = repo_index.get(key) if existing is None: is_new = True elif existing['RIFER'] != new_rifer or existing['COD_IMDB'] != new_imdb: is_new = False # aggiornato else: continue # invariato — non salvare params_list.append(( fornitore or '', titolo, rec.get('TIPO_TITOLO', 'TO'), rec.get('ANNO') or None, rec.get('REGISTA') or None, new_rifer or None, new_imdb or None, now, now, )) if is_new: nuovi += 1 else: aggiornati += 1 if params_list: try: self._db_executemany(_UPSERT_SQL, params_list) except Exception as e: logger.error(f"save_to_repository: {e}") logger.info(f"LINKER_REPOSITORY: {nuovi} nuovi, {aggiornati} aggiornati su {len(records)} candidati") return nuovi, aggiornati # ------------------------------------------------------------------ # # CUSTOM_ANAGR_ORIZ # # ------------------------------------------------------------------ # def get_all_custom(self) -> pd.DataFrame: return self.catalog.get_all_custom() def get_custom_by_rifer(self, rifer: str) -> pd.DataFrame: return self.catalog.get_custom_by_rifer(rifer) def get_custom_by_imdb(self, imdb_code: str) -> pd.DataFrame: return self.catalog.get_custom_by_imdb(imdb_code) def get_custom_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_custom_by_rifers(rifers) # ------------------------------------------------------------------ # # IMDB # # ------------------------------------------------------------------ # def get_all_imdb(self) -> pd.DataFrame: return self.catalog.get_all_imdb() def get_imdb_by_code(self, imdb_code: str) -> pd.DataFrame: return self.catalog.get_imdb_by_code(imdb_code) def search_imdb(self, term: str, by_regista: bool = False) -> pd.DataFrame: return self.catalog.search_imdb(term, by_regista) def search_custom(self, term: str, by_regista: bool = False) -> pd.DataFrame: return self.catalog.search_custom(term, by_regista) # ------------------------------------------------------------------ # # DIRITTI # # ------------------------------------------------------------------ # def get_diritti_full_by_rifer(self, rifer: str) -> pd.DataFrame: return self.catalog.get_diritti_full_by_rifer(rifer) def get_diritti_full_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_diritti_full_by_rifers(rifers) # ------------------------------------------------------------------ # # EMESSO # # ------------------------------------------------------------------ # def get_emesso_both_first_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_emesso_both_first_by_rifers(rifers) def get_emesso_both_last_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_emesso_both_last_by_rifers(rifers) def get_all_imdb_rifer_map(self) -> dict: return self.catalog.get_all_imdb_rifer_map() def get_all_emesso_agg(self): return self.catalog.get_all_emesso_agg() def get_all_gemma(self) -> pd.DataFrame: return self.catalog.get_all_gemma() def get_emesso_last_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_emesso_last_by_rifers(rifers) def get_emesso_tematiche_last_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_emesso_tematiche_last_by_rifers(rifers) def get_gemma_valutazioni_by_imdb_codes(self, imdb_codes: list) -> pd.DataFrame: return self.catalog.get_gemma_valutazioni_by_imdb_codes(imdb_codes) def get_rifer_by_imdb_codes(self, imdb_codes: list) -> pd.DataFrame: return self.catalog.get_rifer_by_imdb_codes(imdb_codes) def get_imdb_data_by_codes(self, imdb_codes: list) -> pd.DataFrame: return self.catalog.get_imdb_data_by_codes(imdb_codes) def get_scelte_rete_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_scelte_rete_by_rifers(rifers) def get_boxoffice_by_rifers(self, rifers: list) -> pd.DataFrame: return self.catalog.get_boxoffice_by_rifers(rifers) def get_emissioni_tematiche_by_rifer(self, rifer: str) -> pd.DataFrame: return self.catalog.get_emissioni_tematiche_by_rifer(rifer) def get_emissioni_generaliste_by_rifer(self, rifer: str) -> pd.DataFrame: return self.catalog.get_emissioni_generaliste_by_rifer(rifer) def get_edizioni_by_rifer(self, rifer: str) -> pd.DataFrame: return self.catalog.get_edizioni_by_rifer(rifer) def get_prodotti_by_superserie(self, superserie_name: str) -> pd.DataFrame: return self.catalog.get_prodotti_by_superserie(superserie_name) # ------------------------------------------------------------------ # # LINKER_VALUTAZIONE # # ------------------------------------------------------------------ # _NO_IMDB_TOKENS = {'no imdb', 'noimdb', 'n/a', '-'} def _imdb_mancante(self, val: str) -> bool: v = val.strip().lower() return not v or v in self._NO_IMDB_TOKENS def _ensure_linker_valutazione_on(self, con: sqlite3.Connection): con.execute(""" CREATE TABLE IF NOT EXISTS LINKER_VALUTAZIONE ( COD_IMDB VARCHAR(255), FORNITORE VARCHAR(255), DATA VARCHAR(255), C5 VARCHAR(255), I1 VARCHAR(255), R4 VARCHAR(255), LA5 VARCHAR(255), I2 VARCHAR(255), IRIS VARCHAR(255), TOPCRIME VARCHAR(255), FOCUS VARCHAR(255), C20 VARCHAR(255), CINE34 VARCHAR(255), TWENTYSEVEN VARCHAR(255), GEMSIN VARCHAR(255), data_creazione VARCHAR(255) ) """) con.commit() # ------------------------------------------------------------------ # # COMPETITIVE_IMDB_MAP # # ------------------------------------------------------------------ # def ensure_competitive_imdb_map(self): """Crea la tabella COMPETITIVE_IMDB_MAP se non esiste.""" self.sqlite_con.execute(""" CREATE TABLE IF NOT EXISTS COMPETITIVE_IMDB_MAP ( TITOLO VARCHAR(255) NOT NULL, DISTRIBUTORE VARCHAR(255) NOT NULL, IMDB_CODE VARCHAR(255), data_update VARCHAR(255), PRIMARY KEY (TITOLO, DISTRIBUTORE) ) """) self.sqlite_con.commit() def get_competitive_imdb_map(self) -> pd.DataFrame: """Restituisce tutti gli agganci Titolo+Distributore→IMDB.""" try: return pd.read_sql_query( "SELECT TITOLO, DISTRIBUTORE, IMDB_CODE FROM COMPETITIVE_IMDB_MAP", self.sqlite_con, ) except Exception as e: logger.error(f"get_competitive_imdb_map: {e}") return pd.DataFrame(columns=['TITOLO', 'DISTRIBUTORE', 'IMDB_CODE']) def upsert_competitive_imdb(self, titolo: str, distributore: str, imdb_code: str): """INSERT OR REPLACE dell'aggancio IMDB per Titolo+Distributore.""" now = datetime.now().isoformat(sep=' ') self._db_execute( """INSERT INTO COMPETITIVE_IMDB_MAP (TITOLO, DISTRIBUTORE, IMDB_CODE, data_update) VALUES (?, ?, ?, ?) ON CONFLICT(TITOLO, DISTRIBUTORE) DO UPDATE SET IMDB_CODE = excluded.IMDB_CODE, data_update = excluded.data_update""", (titolo.strip(), distributore.strip(), imdb_code.strip(), now), ) # ------------------------------------------------------------------ # # OTT_EXT (storico batch — append-only) # # ------------------------------------------------------------------ # def _ensure_ott_ext_on(self, con: sqlite3.Connection): """Applica schema OTT_EXT su una connessione.""" con.execute(""" CREATE TABLE IF NOT EXISTS OTT_EXT ( mediaset_id VARCHAR(255) NOT NULL, provider VARCHAR(255) NOT NULL, data_creazione VARCHAR(255) NOT NULL, segnalazioni_automatiche VARCHAR(255), imdb_override VARCHAR(255), PRIMARY KEY (mediaset_id, provider COLLATE NOCASE) ) """) try: con.execute("ALTER TABLE OTT_EXT ADD COLUMN imdb_override VARCHAR(255)") except Exception: pass existing = {r[1] for r in con.execute("PRAGMA table_info(OTT_EXT)").fetchall()} legacy_cols = {'budget', 'gross_world', 'gross_us_canada', 'production_company', 'distributor'} if legacy_cols & existing: con.execute(""" CREATE TABLE OTT_EXT_new ( mediaset_id VARCHAR(255) NOT NULL, provider VARCHAR(255) NOT NULL, data_creazione VARCHAR(255) NOT NULL, segnalazioni_automatiche VARCHAR(255), imdb_override VARCHAR(255), PRIMARY KEY (mediaset_id, provider COLLATE NOCASE) ) """) con.execute(""" INSERT INTO OTT_EXT_new (mediaset_id, provider, data_creazione, segnalazioni_automatiche, imdb_override) SELECT mediaset_id, provider, data_creazione, segnalazioni_automatiche, imdb_override FROM OTT_EXT """) con.execute("DROP TABLE OTT_EXT") con.execute("ALTER TABLE OTT_EXT_new RENAME TO OTT_EXT") con.commit() def ensure_ott_ext(self): self._ensure_ott_ext_on(self.sqlite_con) if self.server_con is not None: try: self._ensure_ott_ext_on(self.server_con) except Exception as e: logger.warning(f"[OTT server] ensure_ott_ext: {e}") def save_imdb_override(self, mediaset_id: str, provider: str, imdb_code: str) -> bool: """Salva il codice IMDB inserito manualmente in OTT_EXT.""" try: return self._db_execute( "UPDATE OTT_EXT SET imdb_override=? WHERE mediaset_id=? AND provider=?", (imdb_code or None, mediaset_id, provider) ) except Exception as e: logger.error(f"save_imdb_override: {e}") return False def save_ott_ext_batch(self, records: list[dict], data_creazione: str) -> int: """Aggiunge a OTT_EXT le coppie (mediaset_id, provider) mai viste prima. Poi ricalcola segnalazioni_automatiche su TUTTI i mediaset_id in OTT_EXT: - SMONTATO: era su SVOD ma ora ha solo finestre scadute - NULL: ha almeno una finestra SVOD ancora attiva (data_fine vuota o >= oggi) """ from datetime import date today = date.today().isoformat() keys_in_ott = { (rec.get('mediaset_id') or '', rec.get('provider') or '') for rec in records } rows = [(k[0], k[1], data_creazione, None) for k in keys_in_ott] active_ids: set[str] = set() svod_ids: set[str] = set() for rec in records: if str(rec.get('tipo_finestra', '')).upper() == 'SVOD': mid = rec.get('mediaset_id') or '' svod_ids.add(mid) data_fine = (rec.get('data_fine') or '').strip() if not data_fine or data_fine >= today: active_ids.add(mid) expired_ids = svod_ids - active_ids def _apply(con: sqlite3.Connection, label: str): cur = con.executemany( "INSERT OR IGNORE INTO OTT_EXT " "(mediaset_id, provider, data_creazione, segnalazioni_automatiche) " "VALUES (?,?,?,?)", rows, ) inserted = cur.rowcount if cur.rowcount >= 0 else -1 if active_ids: con.executemany( "UPDATE OTT_EXT SET segnalazioni_automatiche = NULL " "WHERE mediaset_id = ? AND segnalazioni_automatiche IS NOT NULL", [(mid,) for mid in active_ids], ) if expired_ids: con.executemany( "UPDATE OTT_EXT SET segnalazioni_automatiche = 'SMONTATO' " "WHERE mediaset_id = ?", [(mid,) for mid in expired_ids], ) con.commit() logger.info( f"[OTT] save_ott_ext_batch [{label}]: {len(rows)} coppie totali, " f"{inserted} nuove inserite, {len(active_ids)} attive, {len(expired_ids)} scadute" ) if self.server_con is not None: srv = self.server_con import threading threading.Thread(target=lambda: _apply(srv, 'server'), daemon=True).start() try: _apply(self.sqlite_con, 'locale') except Exception as e: logger.error(f"save_ott_ext_batch: {e}") self.sqlite_con.rollback() return 0 return len(rows) def recalcola_segnalazioni_smontato(self, records: list[dict]) -> tuple[int, int]: """Ricalcola segnalazioni_automatiche su tutto OTT_EXT dai record OTT attuali. Ritorna (reset_a_null, marcati_smontato). """ from datetime import date today = date.today().isoformat() active_ids: set[str] = set() svod_ids: set[str] = set() for rec in records: if str(rec.get('tipo_finestra', '')).upper() == 'SVOD': mid = rec.get('mediaset_id') or '' svod_ids.add(mid) data_fine = (rec.get('data_fine') or '').strip() if not data_fine or data_fine >= today: active_ids.add(mid) expired_ids = svod_ids - active_ids def _apply(con: sqlite3.Connection) -> tuple[int, int]: r = s = 0 if active_ids: cur = con.executemany( "UPDATE OTT_EXT SET segnalazioni_automatiche = NULL " "WHERE mediaset_id = ? AND segnalazioni_automatiche IS NOT NULL", [(mid,) for mid in active_ids], ) r = cur.rowcount if cur.rowcount >= 0 else len(active_ids) if expired_ids: cur = con.executemany( "UPDATE OTT_EXT SET segnalazioni_automatiche = 'SMONTATO' " "WHERE mediaset_id = ?", [(mid,) for mid in expired_ids], ) s = cur.rowcount if cur.rowcount >= 0 else len(expired_ids) con.commit() return r, s if self.server_con is not None: srv = self.server_con import threading threading.Thread(target=lambda: _apply(srv), daemon=True).start() reset = smontato = 0 try: reset, smontato = _apply(self.sqlite_con) except Exception as e: logger.error(f"recalcola_segnalazioni_smontato: {e}") self.sqlite_con.rollback() return reset, smontato def get_ott_ext_last_date(self) -> str | None: try: row = self.sqlite_con.execute( "SELECT MAX(data_creazione) FROM OTT_EXT" ).fetchone() return row[0] if row and row[0] else None except Exception as e: logger.error(f"get_ott_ext_last_date: {e}") return None def get_ott_ext_cutoff_date(self, n_batches: int = 3) -> str | None: """Restituisce la data cutoff = n-esima data più recente nei batch OTT_EXT.""" try: rows = self.sqlite_con.execute( "SELECT DISTINCT data_creazione FROM OTT_EXT " "ORDER BY data_creazione DESC LIMIT ?", (n_batches,) ).fetchall() if not rows: return None return rows[-1][0] # la più vecchia delle ultime n except Exception as e: logger.error(f"get_ott_ext_cutoff_date: {e}") return None def get_ott_ext_last_dates(self, n_batches: int = 3) -> list[str]: """Restituisce le ultime n date distinte di OTT_EXT, dalla più recente.""" try: rows = self.sqlite_con.execute( "SELECT DISTINCT data_creazione FROM OTT_EXT " "ORDER BY data_creazione DESC LIMIT ?", (n_batches,) ).fetchall() return [r[0] for r in rows if r[0]] except Exception as e: logger.error(f"get_ott_ext_last_dates: {e}") return [] def get_ott_ext_creazione_map(self, since_date: str = None) -> dict: """Restituisce {(mediaset_id, provider): (data_creazione, segnalazioni_automatiche, imdb_override)}. Se since_date è fornita, solo le coppie con data_creazione >= since_date.""" try: if since_date: rows = self.sqlite_con.execute( "SELECT mediaset_id, provider, data_creazione, segnalazioni_automatiche, imdb_override " "FROM OTT_EXT WHERE data_creazione >= ?", (since_date,) ).fetchall() else: rows = self.sqlite_con.execute( "SELECT mediaset_id, provider, data_creazione, segnalazioni_automatiche, imdb_override " "FROM OTT_EXT" ).fetchall() return {(r[0], r[1]): (r[2], r[3], r[4]) for r in rows} except Exception as e: logger.error(f"get_ott_ext_creazione_map: {e}") return {} def get_ott_joined_for_mode(self, mode: str, anno_min: int = 0, cutoff: str = None) -> 'pd.DataFrame': """Delegate a ParquetCatalog: DuckDB su ott.parquet + ott_ext in-memory.""" return self.catalog.get_ott_joined_for_mode(mode, anno_min, cutoff) # ------------------------------------------------------------------ # # tabValutazioniExtraGemma (OTT — valutazioni reti per titolo) # # ------------------------------------------------------------------ # _NETWORKS_VEG = ['C5', 'I1', 'R4', 'LA5', 'I2', 'IRIS', 'TOPCRIME', 'FOCUS', 'C20', 'CINE34', 'TWENTYSEVEN'] def _ensure_valutazioni_extra_gemma_on(self, con: sqlite3.Connection): con.execute(""" CREATE TABLE IF NOT EXISTS tabValutazioniExtraGemma ( IMDB VARCHAR(255) NOT NULL, Ambito VARCHAR(255), C5 VARCHAR(255), I1 VARCHAR(255), R4 VARCHAR(255), LA5 VARCHAR(255), I2 VARCHAR(255), IRIS VARCHAR(255), TOPCRIME VARCHAR(255), FOCUS VARCHAR(255), C20 VARCHAR(255), CINE34 VARCHAR(255), TWENTYSEVEN VARCHAR(255), SintesiPropostaGiro VARCHAR(255), NoIcrMotivazione VARCHAR(255), NoIcr INTEGER DEFAULT 0, data_creazione VARCHAR(255), PRIMARY KEY (IMDB, Ambito) ) """) con.commit() def ensure_valutazioni_extra_gemma(self): self._ensure_valutazioni_extra_gemma_on(self.sqlite_con) if self.server_con is not None: try: self._ensure_valutazioni_extra_gemma_on(self.server_con) except Exception as e: logger.warning(f"[OTT server] ensure_valutazioni_extra_gemma: {e}") def get_valutazioni_extra_gemma(self, imdb: str, ambito: str) -> dict | None: """Cerca per (IMDB, Ambito) esatto; se non trovato, fallback su IMDB solo (primo record).""" try: df = pd.read_sql_query( "SELECT * FROM tabValutazioniExtraGemma WHERE IMDB=? AND Ambito=?", self.sqlite_con, params=(imdb, ambito) ) if not df.empty: return df.iloc[0].to_dict() df2 = pd.read_sql_query( "SELECT * FROM tabValutazioniExtraGemma WHERE IMDB=? LIMIT 1", self.sqlite_con, params=(imdb,) ) return df2.iloc[0].to_dict() if not df2.empty else None except Exception as e: logger.error(f"get_valutazioni_extra_gemma: {e}") return None @staticmethod def _calcola_sintesi_proposta_giro(data: dict) -> str | None: if data.get('NoIcr'): return 'NO ICR' nets = ['C5', 'I1', 'R4', 'LA5', 'I2', 'IRIS', 'TOPCRIME', 'FOCUS', 'C20', 'CINE34', 'TWENTYSEVEN'] filled = [data[n] for n in nets if data.get(n)] if not filled: return None has_bang = any(v == '!' for v in filled) has_other = any(v != '!' for v in filled) if has_bang and has_other: return 'vp' return '!' if has_bang else 'V' def save_valutazioni_extra_gemma(self, imdb: str, ambito: str, data: dict) -> bool: now = datetime.now().isoformat(sep=' ') sintesi = self._calcola_sintesi_proposta_giro(data) sql = """ INSERT INTO tabValutazioniExtraGemma (IMDB, Ambito, C5, I1, R4, LA5, I2, IRIS, TOPCRIME, FOCUS, C20, CINE34, TWENTYSEVEN, SintesiPropostaGiro, NoIcrMotivazione, NoIcr, data_creazione) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(IMDB, Ambito) DO UPDATE SET C5=excluded.C5, I1=excluded.I1, R4=excluded.R4, LA5=excluded.LA5, I2=excluded.I2, IRIS=excluded.IRIS, TOPCRIME=excluded.TOPCRIME, FOCUS=excluded.FOCUS, C20=excluded.C20, CINE34=excluded.CINE34, TWENTYSEVEN=excluded.TWENTYSEVEN, SintesiPropostaGiro=excluded.SintesiPropostaGiro, NoIcrMotivazione=excluded.NoIcrMotivazione, NoIcr=excluded.NoIcr, data_creazione=excluded.data_creazione """ params = ( imdb, ambito, data.get('C5') or None, data.get('I1') or None, data.get('R4') or None, data.get('LA5') or None, data.get('I2') or None, data.get('IRIS') or None, data.get('TOPCRIME') or None, data.get('FOCUS') or None, data.get('C20') or None, data.get('CINE34') or None, data.get('TWENTYSEVEN') or None, sintesi, data.get('NoIcrMotivazione') or None, 1 if data.get('NoIcr') else 0, now, ) try: return self._db_execute(sql, params) except Exception as e: logger.error(f"save_valutazioni_extra_gemma: {e}") return False def delete_valutazioni_extra_gemma(self, imdb: str, ambito: str) -> bool: try: return self._db_execute( "DELETE FROM tabValutazioniExtraGemma WHERE IMDB=? AND Ambito=?", (imdb, ambito) ) except Exception as e: logger.error(f"delete_valutazioni_extra_gemma: {e}") return False def get_valutazioni_extra_gemma_map(self) -> dict: """Ritorna {imdb: SintesiPropostaGiro} — join solo su IMDB, nessun filtro Ambito.""" try: rows = self.sqlite_con.execute( "SELECT IMDB, SintesiPropostaGiro FROM tabValutazioniExtraGemma" ).fetchall() return {row[0]: (row[1] or '') for row in rows if row[0]} except Exception as e: logger.error(f"get_valutazioni_extra_gemma_map: {e}") return {} def get_valutazioni_extra_gemma_full_map(self) -> dict: """Ritorna {imdb: row_dict} con tutti i campi di tabValutazioniExtraGemma. Se esistono più record per lo stesso IMDB (ambiti diversi), usa il primo.""" try: df = pd.read_sql_query( "SELECT * FROM tabValutazioniExtraGemma", self.sqlite_con ) if df.empty: return {} result = {} for _, row in df.iterrows(): imdb = str(row.get('IMDB', '') or '').strip() if imdb and imdb not in result: result[imdb] = row.to_dict() return result except Exception as e: logger.error(f"get_valutazioni_extra_gemma_full_map: {e}") return {} # ------------------------------------------------------------------ # # LINKER_VALUTAZIONE # # ------------------------------------------------------------------ # def get_fornitori_valutazione(self) -> list[str]: """Restituisce la lista distinta di FORNITORE da LINKER_VALUTAZIONE.""" self._ensure_linker_valutazione_on(self.sqlite_con) try: cur = self.sqlite_con.cursor() cur.execute( "SELECT DISTINCT FORNITORE FROM LINKER_VALUTAZIONE " "WHERE FORNITORE IS NOT NULL ORDER BY FORNITORE" ) return [r[0] for r in cur.fetchall()] except Exception as e: logger.error(f"get_fornitori_valutazione: {e}") return [] def salva_valutazioni(self, rows: list[dict], fornitore: str, data: str, chunk_size: int = 50, on_chunk=None) -> tuple[int, int]: """ Inserisce le righe in LINKER_VALUTAZIONE in chunk (dual-write server+locale). on_chunk(saved: int, total: int) — chiamato dopo ogni chunk. Restituisce (salvate, saltate_per_no_imdb). """ _COL_MAP = { 'C5': 'C5', 'I1': 'I1', 'R4': 'R4', 'LA5': 'LA5', 'I2': 'I2', 'IRIS': 'IRIS', 'TOP CRIME': 'TOPCRIME', 'FOCUS': 'FOCUS', 'C20': 'C20', 'CINE34': 'CINE34', 'TWENTY SEVEN': 'TWENTYSEVEN', 'IMDB': 'COD_IMDB', } self._ensure_linker_valutazione_on(self.sqlite_con) if self.server_con is not None: self._ensure_linker_valutazione_on(self.server_con) # Schema fisso — tutte le righe devono avere le stesse colonne per executemany _DB_COLS = ['COD_IMDB', 'FORNITORE', 'DATA', 'C5', 'I1', 'R4', 'LA5', 'I2', 'IRIS', 'TOPCRIME', 'FOCUS', 'C20', 'CINE34', 'TWENTYSEVEN', 'data_creazione'] saltate = 0 now = datetime.now().strftime('%Y-%m-%d %H:%M:%S') params_list = [] for row in rows: if self._imdb_mancante(row.get('IMDB', '')): saltate += 1 continue db_row = {col: '' for col in _DB_COLS} for k, v in row.items(): if k in _COL_MAP: db_row[_COL_MAP[k]] = v db_row['FORNITORE'] = fornitore db_row['DATA'] = data db_row['data_creazione'] = now params_list.append(db_row) salvate = len(params_list) if params_list: cols = ', '.join(_DB_COLS) placeholders = ', '.join(['?'] * len(_DB_COLS)) sql = f"INSERT INTO LINKER_VALUTAZIONE ({cols}) VALUES ({placeholders})" written = 0 for i in range(0, salvate, chunk_size): chunk = params_list[i:i + chunk_size] self._db_executemany(sql, [list(p.values()) for p in chunk]) written += len(chunk) if on_chunk: on_chunk(written, salvate) # Cancella gemma_unified.db e lancia subito la ricostruzione in background gemma_unified = self.catalog.parquet_path / 'gemma_unified.db' try: if gemma_unified.exists(): gemma_unified.unlink() except Exception: pass self._rebuild_gemma_unified_async() return salvate, saltate def _rebuild_gemma_unified_async(self): """Lancia ensure_gemma_unified in un thread daemon usando linker.db locale.""" try: local_db_path = None rows = self.sqlite_con.execute("PRAGMA database_list").fetchall() for _, name, filename in rows: if name == 'main' and filename: local_db_path = Path(filename) break if local_db_path is None: return parquet_dir = self.catalog.parquet_path import threading as _th from launcher_src.gemma_unified_builder import ensure_gemma_unified _th.Thread( target=ensure_gemma_unified, args=(parquet_dir, local_db_path), daemon=True, ).start() except Exception as e: logger.warning(f"_rebuild_gemma_unified_async: {e}") # ------------------------------------------------------------------ # # PRODUCT ANNOTATIONS — NOTA ICR / NETWORK FREE / FORNITORE DIRITTI # # ------------------------------------------------------------------ # def _ensure_product_annotations_on(self, con: sqlite3.Connection): con.execute(""" CREATE TABLE IF NOT EXISTS product_annotations ( rifer INTEGER, imdb_code VARCHAR(255), ott_mediaset_id VARCHAR(255), nota_icr VARCHAR(255), network_free VARCHAR(255), fornitore_diritti VARCHAR(255), updated_at VARCHAR(30), CHECK (rifer IS NOT NULL OR imdb_code IS NOT NULL OR ott_mediaset_id IS NOT NULL) ) """) # Migration: add ott_mediaset_id if missing (existing DBs) existing_cols = {r[1] for r in con.execute("PRAGMA table_info(product_annotations)").fetchall()} if 'ott_mediaset_id' not in existing_cols: con.execute("ALTER TABLE product_annotations ADD COLUMN ott_mediaset_id VARCHAR(255)") con.execute( "CREATE UNIQUE INDEX IF NOT EXISTS idx_pa_rifer ON product_annotations (rifer) WHERE rifer IS NOT NULL" ) con.execute( "CREATE UNIQUE INDEX IF NOT EXISTS idx_pa_imdb ON product_annotations (imdb_code) WHERE imdb_code IS NOT NULL AND rifer IS NULL" ) con.execute( "CREATE UNIQUE INDEX IF NOT EXISTS idx_pa_ott ON product_annotations (ott_mediaset_id) " "WHERE ott_mediaset_id IS NOT NULL AND imdb_code IS NULL AND rifer IS NULL" ) # OTT rows have rifer=NULL and imdb_code=NULL — disable the legacy CHECK constraint con.execute("PRAGMA ignore_check_constraints = 1") con.commit() def ensure_product_annotations(self): self._ensure_product_annotations_on(self.sqlite_con) if self.server_con is not None: try: self._ensure_product_annotations_on(self.server_con) except Exception as e: logger.warning(f"[OTT server] ensure_product_annotations: {e}") def get_product_annotation(self, rifer=None, imdb_code: str = None) -> dict | None: """Lookup per RIFER (priorità) poi imdb_code. Ritorna dict o None.""" try: if rifer: row = self.sqlite_con.execute( "SELECT nota_icr, network_free, fornitore_diritti FROM product_annotations WHERE rifer=?", (int(rifer),) ).fetchone() if row: return {'nota_icr': row[0], 'network_free': row[1], 'fornitore_diritti': row[2]} if imdb_code: row = self.sqlite_con.execute( "SELECT nota_icr, network_free, fornitore_diritti FROM product_annotations WHERE imdb_code=?", (imdb_code,) ).fetchone() if row: return {'nota_icr': row[0], 'network_free': row[1], 'fornitore_diritti': row[2]} except Exception as e: logger.error(f"get_product_annotation: {e}") return None def save_product_annotation(self, rifer=None, imdb_code: str = None, ott_mediaset_id: str = None, nota_icr: str = None, network_free: str = None, fornitore_diritti: str = None) -> bool: if not rifer and not imdb_code and not ott_mediaset_id: return False now = datetime.now().isoformat(sep=' ') try: if rifer: return self._db_execute(""" INSERT INTO product_annotations (rifer, imdb_code, ott_mediaset_id, nota_icr, network_free, fornitore_diritti, updated_at) VALUES (?, ?, NULL, ?, ?, ?, ?) ON CONFLICT(rifer) WHERE rifer IS NOT NULL DO UPDATE SET imdb_code=excluded.imdb_code, nota_icr=excluded.nota_icr, network_free=excluded.network_free, fornitore_diritti=excluded.fornitore_diritti, updated_at=excluded.updated_at """, (int(rifer), imdb_code or None, nota_icr or None, network_free or None, fornitore_diritti or None, now)) elif imdb_code: return self._db_execute(""" INSERT INTO product_annotations (rifer, imdb_code, ott_mediaset_id, nota_icr, network_free, fornitore_diritti, updated_at) VALUES (NULL, ?, NULL, ?, ?, ?, ?) ON CONFLICT(imdb_code) WHERE imdb_code IS NOT NULL AND rifer IS NULL DO UPDATE SET nota_icr=excluded.nota_icr, network_free=excluded.network_free, fornitore_diritti=excluded.fornitore_diritti, updated_at=excluded.updated_at """, (imdb_code, nota_icr or None, network_free or None, fornitore_diritti or None, now)) else: return self._db_execute(""" INSERT INTO product_annotations (rifer, imdb_code, ott_mediaset_id, nota_icr, network_free, fornitore_diritti, updated_at) VALUES (NULL, NULL, ?, ?, ?, ?, ?) ON CONFLICT(ott_mediaset_id) WHERE ott_mediaset_id IS NOT NULL AND imdb_code IS NULL AND rifer IS NULL DO UPDATE SET nota_icr=excluded.nota_icr, network_free=excluded.network_free, fornitore_diritti=excluded.fornitore_diritti, updated_at=excluded.updated_at """, (ott_mediaset_id, nota_icr or None, network_free or None, fornitore_diritti or None, now)) except Exception as e: logger.error(f"save_product_annotation: {e}") return False def get_product_annotations_map(self, rifers: list = None) -> dict: """Ritorna {rifer: {nota_icr, network_free, fornitore_diritti}} per i RIFER forniti. Se rifers=None carica tutto (per bulk join).""" try: if rifers is not None: if not rifers: return {} # Chunked per evitare limite 999 bind variables result = {} chunk_size = 900 for i in range(0, len(rifers), chunk_size): chunk = rifers[i:i + chunk_size] placeholders = ','.join(['?' for _ in chunk]) rows = self.sqlite_con.execute( f"SELECT rifer, nota_icr, network_free, fornitore_diritti FROM product_annotations WHERE rifer IN ({placeholders})", [int(r) for r in chunk] ).fetchall() for row in rows: result[row[0]] = {'nota_icr': row[1], 'network_free': row[2], 'fornitore_diritti': row[3]} return result else: rows = self.sqlite_con.execute( "SELECT rifer, nota_icr, network_free, fornitore_diritti FROM product_annotations WHERE rifer IS NOT NULL" ).fetchall() return {row[0]: {'nota_icr': row[1], 'network_free': row[2], 'fornitore_diritti': row[3]} for row in rows} except Exception as e: logger.error(f"get_product_annotations_map: {e}") return {} def get_product_annotations_imdb_map(self) -> dict: """Ritorna {imdb_code: {nota_icr, network_free, fornitore_diritti}} — solo record con imdb_code.""" try: rows = self.sqlite_con.execute( "SELECT imdb_code, nota_icr, network_free, fornitore_diritti " "FROM product_annotations WHERE imdb_code IS NOT NULL" ).fetchall() return {row[0]: {'nota_icr': row[1], 'network_free': row[2], 'fornitore_diritti': row[3]} for row in rows} except Exception as e: logger.error(f"get_product_annotations_imdb_map: {e}") return {} def get_product_annotations_ott_map(self) -> dict: """Ritorna {ott_mediaset_id: {nota_icr, network_free, fornitore_diritti}} — solo record OTT.""" try: rows = self.sqlite_con.execute( "SELECT ott_mediaset_id, nota_icr, network_free, fornitore_diritti " "FROM product_annotations WHERE ott_mediaset_id IS NOT NULL AND imdb_code IS NULL AND rifer IS NULL" ).fetchall() return {row[0]: {'nota_icr': row[1], 'network_free': row[2], 'fornitore_diritti': row[3]} for row in rows} except Exception as e: logger.error(f"get_product_annotations_ott_map: {e}") return {}