# C:\PYTHON\MyICR_Suite\app\features\dash_diritti\data_loaders.py """ Data Loaders for Dashboard Diritti - Refactored to use CancelableThread """ import pandas as pd from PyQt6.QtCore import pyqtSignal from loguru import logger from app.utils.cancelable_thread import CancelableThread from app.utils.cache_manager import get_cache_manager from .queries import get_diritti_pivot_data, get_diritti_detail_data class DirittiDataLoader(CancelableThread): """ Worker thread per caricare ed elaborare i dati per la pivot dei diritti. Refactored to use CancelableThread for safe cancellation. """ data_loaded = pyqtSignal(pd.DataFrame) def __init__(self, filters: dict): super().__init__() self.filters = filters def _calculate_days_in_year(self, row, year): """Calcola i giorni di validità di un diritto all'interno di un anno specifico.""" start_of_year = pd.Timestamp(year, 1, 1) end_of_year = pd.Timestamp(year, 12, 31) # Trova l'intervallo di sovrapposizione tra il diritto e l'anno overlap_start = max(row['decr'], start_of_year) overlap_end = min(row['scad'], end_of_year) # Calcola la durata in giorni if overlap_start > overlap_end: return 0 return (overlap_end - overlap_start).days + 1 def do_work(self): """ Carica ed elabora i dati per la pivot dei diritti. Override di CancelableThread.do_work() """ if self.is_cancelled: logger.info("DirittiDataLoader: Cancellato prima dell'avvio") return None try: logger.info("Worker Diritti: Esecuzione query dati grezzi...") # Usa CacheManager per cache della query raw data # Nota: cachamo solo la query, non l'elaborazione che dipende dai filtri dinamici cache = get_cache_manager() raw_data = cache.cached_query( namespace='diritti_raw', params=self.filters, query_func=lambda: get_diritti_pivot_data(self.filters), ttl=600 # 10 minuti ) if self.is_cancelled: logger.info("DirittiDataLoader: Cancellato dopo query") return None if raw_data is None or raw_data.empty: self.data_loaded.emit(pd.DataFrame()) return pd.DataFrame() logger.info(f"Worker Diritti: Inizio elaborazione dati ({len(raw_data)} righe).") raw_data['decr'] = pd.to_datetime(raw_data['decr']) raw_data['scad'] = pd.to_datetime(raw_data['scad']) # 1. "Esplode" i diritti, creando una riga per ogni anno di validità exploded_rows = [] filter_start_year = pd.to_datetime(self.filters['data_decorrenza']).year filter_end_year = pd.to_datetime(self.filters['data_scadenza']).year view_mode_decorrenza = self.filters.get('only_decorrenza_year', False) if view_mode_decorrenza: for _, row in raw_data.iterrows(): # Check cancellation periodicamente if self.is_cancelled: logger.info("DirittiDataLoader: Cancellato durante elaborazione") return None decorrenza_year = row['decr'].year if filter_start_year <= decorrenza_year <= filter_end_year: days_in_year = self._calculate_days_in_year(row, decorrenza_year) exploded_rows.append({ 'tipologia': row['tipologia'], 'anno': decorrenza_year, 'fr_rr': row['fr_rr'], 'prod': row['prod'], 'perc': row['perc'], 'days': days_in_year }) else: for _, row in raw_data.iterrows(): # Check cancellation periodicamente if self.is_cancelled: logger.info("DirittiDataLoader: Cancellato durante elaborazione") return None start = max(row['decr'].year, filter_start_year) end = min(row['scad'].year, filter_end_year) for year in range(start, end + 1): days_in_year = self._calculate_days_in_year(row, year) exploded_rows.append({ 'tipologia': row['tipologia'], 'anno': year, 'fr_rr': row['fr_rr'], 'prod': row['prod'], 'perc': row['perc'], 'days': days_in_year }) if not exploded_rows: self.data_loaded.emit(pd.DataFrame()) return pd.DataFrame() exploded_df = pd.DataFrame(exploded_rows) # 2. Applica il filtro sulla disponibilità (se attivo) if self.filters.get('exclude_disponibilita_breve'): if self.is_cancelled: return None logger.info("Worker Diritti: Applicazione filtro disponibilità >= 30 giorni.") # Calcola la disponibilità totale per prodotto per anno days_sum_df = exploded_df.groupby(['anno', 'prod'])['days'].sum().reset_index() # Identifica i prodotti con disponibilità sufficiente valid_prods = days_sum_df[days_sum_df['days'] >= 30][['anno', 'prod']] # Filtra per mantenere solo i prodotti validi exploded_df = pd.merge(exploded_df, valid_prods, on=['anno', 'prod'], how='inner') if exploded_df.empty: self.data_loaded.emit(pd.DataFrame()) return pd.DataFrame() # 3. Applica il filtro percentuale (se attivo) if self.filters.get('exclude_perc_minoritaria'): if self.is_cancelled: return None logger.info("Worker Diritti: Applicazione filtro titolarità >= 100%.") perc_sum_df = exploded_df.groupby(['anno', 'prod'])['perc'].sum().reset_index() full_rights_prods = perc_sum_df[perc_sum_df['perc'] >= 100.0][['anno', 'prod']] exploded_df = pd.merge(exploded_df, full_rights_prods, on=['anno', 'prod'], how='inner') if exploded_df.empty: self.data_loaded.emit(pd.DataFrame()) return pd.DataFrame() # 4. Finalizza i dati per la pivot table if self.is_cancelled: logger.info("DirittiDataLoader: Cancellato prima della finalizzazione") return None logger.info("Worker Diritti: Finalizzazione dati per la pivot...") counts_df = exploded_df.groupby(['tipologia', 'anno', 'fr_rr']).size().reset_index(name='count') final_data_for_pivot = counts_df.pivot_table( index=['tipologia', 'anno'], columns='fr_rr', values='count', aggfunc='sum' ).fillna(0).reset_index() final_data_for_pivot.rename(columns={'F': 'first_run_count', 'R': 're_run_count'}, inplace=True) if 'first_run_count' not in final_data_for_pivot.columns: final_data_for_pivot['first_run_count'] = 0 if 're_run_count' not in final_data_for_pivot.columns: final_data_for_pivot['re_run_count'] = 0 logger.info("Worker Diritti: Elaborazione completata.") # Check cancellation prima di emettere if not self.is_cancelled: self.data_loaded.emit(final_data_for_pivot) return final_data_for_pivot except Exception as e: logger.error(f"Worker Diritti: Errore durante il caricamento/elaborazione dati: {e}", exc_info=True) self.error_occurred.emit(str(e)) return None class DirittiDetailDataLoader(CancelableThread): """ Thread per caricamento dati di dettaglio dei diritti. Gestisce il caricamento asincrono dei dettagli per una cella selezionata nella pivot table, evitando di bloccare l'interfaccia. """ detail_loaded = pyqtSignal(pd.DataFrame) def __init__(self, tipologia: str, anno: str, fr_rr: str, filters: dict): super().__init__() self.tipologia = tipologia self.anno = anno self.fr_rr = fr_rr self.filters = filters def do_work(self): """ Carica i dati di dettaglio con caching. Override di CancelableThread.do_work() """ if self.is_cancelled: logger.info("DirittiDetailDataLoader: Cancellato prima dell'avvio") return None try: logger.info(f"DirittiDetailDataLoader: Caricamento dettaglio per {self.tipologia}, {self.anno}, {self.fr_rr}") # Usa CacheManager per cache delle query cache = get_cache_manager() detail_params = { 'tipologia': self.tipologia, 'anno': self.anno, 'fr_rr': self.fr_rr, 'filters': str(self.filters) # Convert dict to string for cache key } detail_df = cache.cached_query( namespace='diritti_detail', params=detail_params, query_func=lambda: get_diritti_detail_data( self.tipologia, self.anno, self.fr_rr, self.filters ), ttl=600 # 10 minuti ) # Check cancellation prima di emettere if not self.is_cancelled: self.detail_loaded.emit(detail_df) logger.info(f"DirittiDetailDataLoader: Dettaglio caricato, {len(detail_df)} righe") return detail_df except Exception as e: logger.error(f"DirittiDetailDataLoader: Errore durante il caricamento: {e}", exc_info=True) self.error_occurred.emit(str(e)) return None