""" Sincronizzazione shared_db/mediatrack.db tra locale e server. - pull_from_server(): copia server→locale se server più recente (mtime) - push_to_server(): copia locale→server via SQLite backup API + file lock - start_background_sync() / stop_background_sync(): timer ogni N minuti - pull_ott_db(): copia ott.db server→locale (one-way, no merge) """ import sqlite3 import threading import time from pathlib import Path from loguru import logger # Import lazy per evitare circular import all'avvio def _get_sync_config(): from app.config_loader import get_sync_config, get_db_path return get_sync_config(), get_db_path('mediatrack') # --- Timer background --- _timer: threading.Timer = None _timer_lock = threading.Lock() def _pull_cache_path() -> Path: """Percorso del file cache che ricorda l'ultima versione del server sincronizzata.""" from app.config_loader import get_app_root return get_app_root() / "launcher_cache" / "shared_db_pull.json" def _read_pull_cache() -> dict: cache = _pull_cache_path() if cache.exists(): try: import json return json.loads(cache.read_text(encoding="utf-8")) except Exception: pass return {} def _write_pull_cache(mtime: float, size: int): import json cache = _pull_cache_path() cache.parent.mkdir(parents=True, exist_ok=True) cache.write_text(json.dumps({"mtime": mtime, "size": size}), encoding="utf-8") def pull_and_preserve_local() -> bool: """ Sync principale al bootstrap. Principio: server = fonte di verità. 1. Rileva se il server è cambiato dall'ultimo pull. 2. Fa uno snapshot dei record locali. 3. Sovrascrive il locale con il server (server vince). 4. Ri-applica on top i record che esistevano solo in locale (prodotti, eventi, distributori non ancora sul server). 5. Se ci sono aggiunte locali, le pusha al server. Questo garantisce che nessun lavoro locale venga perso anche se nel frattempo un altro utente ha modificato il server. """ try: sync_config, local_db = _get_sync_config() server_path = Path(sync_config.get("server_path", "")) if not server_path.exists(): logger.warning(f"[sync] Server non raggiungibile: {server_path}") return False st = server_path.stat() cached = _read_pull_cache() server_unchanged = ( cached.get("mtime") == st.st_mtime and cached.get("size") == st.st_size and local_db.exists() ) if server_unchanged: logger.debug("[sync] pull_and_preserve_local: server invariato — nessuna azione") return False # ── Snapshot locale pre-pull ────────────────────────────────────── local_snapshot = {} if local_db.exists(): try: snap = sqlite3.connect(str(local_db)) snap.row_factory = sqlite3.Row media_cols = [c[1] for c in snap.execute("PRAGMA table_info(AnagraficaMedia)").fetchall()] log_cols = [c[1] for c in snap.execute("PRAGMA table_info(LogAggiornamenti)").fetchall()] distr_cols = [c[1] for c in snap.execute("PRAGMA table_info(Distributori)").fetchall()] local_snapshot = { "media": snap.execute("SELECT * FROM AnagraficaMedia").fetchall(), "log": snap.execute("SELECT * FROM LogAggiornamenti").fetchall(), "distr": snap.execute("SELECT * FROM Distributori").fetchall(), "media_cols": media_cols, "log_cols": log_cols, "distr_cols": distr_cols, } snap.close() except Exception as e: logger.warning(f"[sync] Snapshot locale fallito: {e}") # ── Pull: server → locale ───────────────────────────────────────── import shutil local_db.parent.mkdir(parents=True, exist_ok=True) tmp = local_db.with_suffix(".tmp") shutil.copy2(str(server_path), str(tmp)) tmp.replace(local_db) _write_pull_cache(st.st_mtime, st.st_size) logger.info("[sync] Pull completato (server → locale)") if not local_snapshot: return True # ── Merge: ri-applica record locali non presenti sul server ─────── con = sqlite3.connect(str(local_db)) media_cols = local_snapshot["media_cols"] log_cols = local_snapshot["log_cols"] distr_cols = local_snapshot["distr_cols"] log_cols_no_id = [c for c in log_cols if c != "id_aggiornamento"] col = {c: i for i, c in enumerate(log_cols)} server_media_ids = set( r[0] for r in con.execute("SELECT id_media FROM AnagraficaMedia").fetchall() ) server_deleted_ids = set( r[0] for r in con.execute( "SELECT id_media FROM AnagraficaMedia WHERE deleted_at IS NOT NULL" ).fetchall() ) server_log_keys = set( (r[0], str(r[1] or ""), str(r[2] or ""), str(r[3] or ""), str(r[4] or ""), str(r[5] or "")) for r in con.execute( "SELECT id_media_fk, id_tipo_evento_fk, " "COALESCE(data_evento,''), COALESCE(dettaglio_contesto,''), COALESCE(note,'')," "COALESCE(id_distributore_fk,'')" " FROM LogAggiornamenti" ).fetchall() ) server_distr_names = set( (r[0] or "").lower() for r in con.execute("SELECT nome FROM Distributori").fetchall() ) ph_m = ",".join(["?"] * len(media_cols)) ph_l = ",".join(["?"] * len(log_cols_no_id)) ph_d = ",".join(["?"] * len(distr_cols)) # Indice G_: (titolo lower, distributore lower) → id_media # Usato per bloccare il re-merge di record MT che duplicano un G_ sul server. server_g_index = { (str(r[0] or "").lower().strip(), str(r[1] or "").lower().strip()): r[2] for r in con.execute( "SELECT titolo_ufficiale, distributore_nome, id_media " "FROM AnagraficaMedia WHERE id_media LIKE 'G_%'" ).fetchall() } # id_media MT locali che hanno un G_ corrispondente sul server → non vanno re-mergiati mt_superseded_by_g = set() m_idx = {c: i for i, c in enumerate(media_cols)} for row in local_snapshot["media"]: id_m = row[0] if not str(id_m).startswith("G_"): key = ( str(row[m_idx["titolo_ufficiale"]] or "").lower().strip(), str(row[m_idx["distributore_nome"]] or "").lower().strip(), ) g_id = server_g_index.get(key) if g_id: mt_superseded_by_g.add(id_m) if mt_superseded_by_g: logger.info( f"[sync] Merge: ignorati {len(mt_superseded_by_g)} record MT " f"sostituiti da G_ sul server: {mt_superseded_by_g}" ) ins_media = ins_log = ins_distr = 0 for row in local_snapshot["media"]: if row[0] not in server_media_ids and row[0] not in mt_superseded_by_g and row[0] not in server_deleted_ids: con.execute( f"INSERT OR IGNORE INTO AnagraficaMedia ({','.join(media_cols)}) VALUES ({ph_m})", tuple(row) ) ins_media += 1 for row in local_snapshot["log"]: id_media_fk = row[col["id_media_fk"]] # Se questo evento appartiene a un MT sostituito, re-linka al G_ corrispondente if id_media_fk in mt_superseded_by_g: mt_row = next( (r for r in local_snapshot["media"] if r[0] == id_media_fk), None ) if mt_row: key = ( str(mt_row[m_idx["titolo_ufficiale"]] or "").lower().strip(), str(mt_row[m_idx["distributore_nome"]] or "").lower().strip(), ) id_media_fk = server_g_index.get(key, id_media_fk) key = ( id_media_fk, str(row[col["id_tipo_evento_fk"]] or ""), str(row[col["data_evento"]] or ""), str(row[col["dettaglio_contesto"]] or ""), str(row[col["note"]] or ""), str(row[col["id_distributore_fk"]] or ""), ) if key not in server_log_keys: vals = list(row[col[c]] for c in log_cols_no_id) vals[log_cols_no_id.index("id_media_fk")] = id_media_fk con.execute( f"INSERT INTO LogAggiornamenti ({','.join(log_cols_no_id)}) VALUES ({ph_l})", tuple(vals) ) ins_log += 1 for row in local_snapshot["distr"]: nome = (row[1] or "").lower() if nome and nome not in server_distr_names: con.execute( f"INSERT OR IGNORE INTO Distributori ({','.join(distr_cols)}) VALUES ({ph_d})", tuple(row) ) ins_distr += 1 # Azzera flag offline dopo merge (dati ora al sicuro) con.execute("UPDATE AnagraficaMedia SET is_offline_insert = 0 WHERE is_offline_insert = 1") con.commit() con.close() if ins_media + ins_log + ins_distr > 0: logger.info( f"[sync] Merge locale: +{ins_media} media, +{ins_log} log, +{ins_distr} distributori" ) push_to_server() return True except Exception as e: logger.error(f"[sync] pull_and_preserve_local fallito: {e}") return False def pull_from_server() -> bool: """ Copia server→locale se il server ha una versione diversa dall'ultima sincronizzata. Confronta mtime+size del server con la cache locale — immune da clock skew client/server. Returns True se la copia è avvenuta, False altrimenti. """ try: sync_config, local_db = _get_sync_config() server_path = Path(sync_config.get("server_path", "")) if not server_path.exists(): logger.warning(f"[sync] Server non raggiungibile: {server_path}") return False st = server_path.stat() server_mtime = st.st_mtime server_size = st.st_size # Confronta con l'ultima versione che abbiamo sincronizzato (non con mtime locale) cached = _read_pull_cache() if (cached.get("mtime") == server_mtime and cached.get("size") == server_size and local_db.exists()): logger.debug("[sync] Pull non necessario — server invariato dall'ultimo sync") return False logger.info(f"[sync] Pull: server → locale ({server_path.name})") local_db.parent.mkdir(parents=True, exist_ok=True) # shutil.copy2 è molto più veloce di sqlite3.backup() su share SMB: # copia il file in un'unica operazione bulk invece di pagina per pagina. # Sicuro al pull: avviene all'avvio, nessun client sta scrivendo sul server. import shutil tmp = local_db.with_suffix(".tmp") shutil.copy2(str(server_path), str(tmp)) tmp.replace(local_db) # Salva in cache mtime+size del server appena sincronizzato _write_pull_cache(server_mtime, server_size) logger.info("[sync] Pull completato") return True except Exception as e: logger.error(f"[sync] Pull fallito: {e}") return False def _ensure_updated_cols_server(con: sqlite3.Connection) -> None: """Migra il DB server aggiungendo colonne mancanti ad AnagraficaMedia.""" try: cols = {r[1] for r in con.execute("PRAGMA table_info(AnagraficaMedia)").fetchall()} for col in ('updated_at', 'updated_by', 'deleted_at'): if col not in cols: con.execute(f"ALTER TABLE AnagraficaMedia ADD COLUMN {col} TEXT") con.commit() except Exception: pass def push_to_server() -> bool: """ Push incrementale: invia al server solo i record non ancora presenti. Non sovrascrive mai record esistenti — server e' fonte di verita'. Logica: - AnagraficaMedia: push per id_media non presente su server - Distributori: push per nome non presente su server (senza id locale, server assegna nuovo id); rimappa FK in locale se id differisce - LogAggiornamenti: push per chiave contenuto (con id_distributore_fk rimappato) """ try: sync_config, local_db = _get_sync_config() if not sync_config.get("push_enabled", True): logger.debug("[sync] Push disabilitato (push_enabled: false)") return False server_path = Path(sync_config.get("server_path", "")) if not server_path.parent.exists(): logger.warning(f"[sync] Server non raggiungibile: {server_path.parent}") return False if not local_db.exists(): logger.warning(f"[sync] DB locale non trovato: {local_db}") return False lock_file = server_path.parent / "mediatrack.lock" LOCK_STALE_SECONDS = 60 if lock_file.exists(): lock_age = time.time() - lock_file.stat().st_mtime if lock_age > LOCK_STALE_SECONDS: logger.warning(f"[sync] Lock stale rilevato ({lock_age:.0f}s) — rimosso") lock_file.unlink(missing_ok=True) waited = 0 while lock_file.exists() and waited < 10: time.sleep(1) waited += 1 if lock_file.exists(): logger.warning("[sync] Push annullato — lock non rilasciato entro 10s") return False lock_file.touch() try: local_con = sqlite3.connect(str(local_db)) local_con.row_factory = sqlite3.Row server_con = sqlite3.connect(str(server_path), timeout=30.0) server_con.row_factory = sqlite3.Row _ensure_updated_cols_server(server_con) try: # ── Schema ────────────────────────────────────────────────── media_cols = [c[1] for c in local_con.execute("PRAGMA table_info(AnagraficaMedia)").fetchall()] log_cols = [c[1] for c in local_con.execute("PRAGMA table_info(LogAggiornamenti)").fetchall()] distr_cols = [c[1] for c in local_con.execute("PRAGMA table_info(Distributori)").fetchall()] log_cols_no_id = [c for c in log_cols if c != "id_aggiornamento"] distr_cols_no_id = [c for c in distr_cols if c != "id_distributore"] ph_m = ",".join(["?"] * len(media_cols)) ph_l = ",".join(["?"] * len(log_cols_no_id)) ph_d = ",".join(["?"] * len(distr_cols_no_id)) # ── Stato server ───────────────────────────────────────────── server_media_rows = { r["id_media"]: r for r in server_con.execute("SELECT * FROM AnagraficaMedia").fetchall() } server_media_ids = set(server_media_rows.keys()) server_log_keys = set( (str(r[0] or ""), str(r[1] or ""), str(r[2] or ""), str(r[3] or ""), str(r[4] or ""), str(r[5] or "")) for r in server_con.execute( "SELECT id_media_fk, id_tipo_evento_fk, " "COALESCE(data_evento,''), COALESCE(dettaglio_contesto,''), " "COALESCE(note,''), COALESCE(id_distributore_fk,'') " "FROM LogAggiornamenti" ).fetchall() ) server_distr_by_nome = { (r[0] or "").lower(): r[1] for r in server_con.execute("SELECT nome, id_distributore FROM Distributori").fetchall() } ins_media = upd_media = ins_log = ins_distr = 0 id_remap = {} # local id_distributore -> server id_distributore # ── 1. Distributori (per nome, senza id locale) ────────────── for row in local_con.execute("SELECT * FROM Distributori").fetchall(): nome = (row["nome"] or "").lower() if not nome: continue if nome in server_distr_by_nome: server_id = server_distr_by_nome[nome] if row["id_distributore"] != server_id: id_remap[row["id_distributore"]] = server_id else: vals = [row[c] for c in distr_cols_no_id] cur = server_con.execute( f"INSERT INTO Distributori ({','.join(distr_cols_no_id)}) VALUES ({ph_d})", tuple(vals) ) new_id = cur.lastrowid id_remap[row["id_distributore"]] = new_id server_distr_by_nome[nome] = new_id ins_distr += 1 # ── 2. AnagraficaMedia ─────────────────────────────────────── # ID soft-deleted sul server: non vanno re-inseriti dai client server_deleted_ids = { r[0] for r in server_con.execute( "SELECT id_media FROM AnagraficaMedia WHERE deleted_at IS NOT NULL" ).fetchall() } # Columns eligible for UPDATE (exclude PK and client-only flag) media_upd_cols = [c for c in media_cols if c not in ("id_media", "is_offline_insert")] ph_upd = ",".join([f"{c}=?" for c in media_upd_cols]) for row in local_con.execute("SELECT * FROM AnagraficaMedia").fetchall(): if row["id_media"] not in server_media_ids: # Non re-inserire record che il server ha soft-deleted if row["id_media"] in server_deleted_ids: continue server_con.execute( f"INSERT OR IGNORE INTO AnagraficaMedia ({','.join(media_cols)}) VALUES ({ph_m})", tuple(row) ) ins_media += 1 else: # Update server only if local record is newer (last-write-wins per record) srv = server_media_rows[row["id_media"]] local_upd = row["updated_at"] if "updated_at" in row.keys() else None srv_upd = srv["updated_at"] if "updated_at" in srv.keys() else None if local_upd and (not srv_upd or local_upd > srv_upd): vals = [row[c] for c in media_upd_cols] + [row["id_media"]] server_con.execute( f"UPDATE AnagraficaMedia SET {ph_upd} WHERE id_media = ?", tuple(vals) ) upd_media += 1 # ── 3. LogAggiornamenti (con remap id distributore) ────────── for row in local_con.execute("SELECT * FROM LogAggiornamenti").fetchall(): id_distr_fk = row["id_distributore_fk"] if id_distr_fk in id_remap: id_distr_fk = id_remap[id_distr_fk] id_media_fk = row["id_media_fk"] key = ( str(id_media_fk or ""), str(row["id_tipo_evento_fk"] or ""), str(row["data_evento"] or ""), str(row["dettaglio_contesto"] or ""), str(row["note"] or ""), str(id_distr_fk or ""), ) if key not in server_log_keys: vals = [ id_distr_fk if c == "id_distributore_fk" else row[c] for c in log_cols_no_id ] server_con.execute( f"INSERT INTO LogAggiornamenti ({','.join(log_cols_no_id)}) VALUES ({ph_l})", tuple(vals) ) ins_log += 1 server_con.commit() # ── 4. Rimappa id distributori nel DB locale ───────────────── # Aggiorna prima LogAggiornamenti (FK), poi Distributori (PK) if id_remap: for old_id, new_id in id_remap.items(): local_con.execute( "UPDATE LogAggiornamenti SET id_distributore_fk = ? " "WHERE id_distributore_fk = ?", (new_id, old_id) ) for old_id, new_id in id_remap.items(): local_con.execute( "UPDATE Distributori SET id_distributore = ? " "WHERE id_distributore = ?", (new_id, old_id) ) local_con.commit() logger.info(f"[sync] Push: rimappati {len(id_remap)} id distributori in locale") if ins_media + upd_media + ins_log + ins_distr > 0: logger.info( f"[sync] Push incrementale: +{ins_media} media nuovi, " f"~{upd_media} media aggiornati, " f"+{ins_log} log, +{ins_distr} distributori" ) # Aggiorna cache: evita pull ridondante al prossimo avvio st = server_path.stat() _write_pull_cache(st.st_mtime, st.st_size) else: logger.debug("[sync] Push: nessuna novita' da sincronizzare") return True finally: local_con.close() server_con.close() finally: lock_file.unlink(missing_ok=True) except Exception as e: logger.error(f"[sync] Push incrementale fallito: {e}") return False def push_offline_inserts() -> int: """ Inserisce nel DB server i prodotti creati offline (is_offline_insert=1). Merge selettivo: non sovrascrive i dati esistenti sul server. Returns: numero di prodotti sincronizzati (0 se nessuno o errore). """ try: sync_config, local_db = _get_sync_config() server_path = Path(sync_config.get("server_path", "")) if not server_path.exists() or not local_db.exists(): return 0 local_conn = sqlite3.connect(str(local_db)) local_conn.row_factory = sqlite3.Row pending = local_conn.execute( "SELECT id_media FROM AnagraficaMedia WHERE is_offline_insert = 1" ).fetchall() if not pending: local_conn.close() return 0 pending_ids = [r["id_media"] for r in pending] logger.info(f"[sync] push_offline_inserts: {len(pending_ids)} prodotti da sincronizzare") # Colonne AnagraficaMedia (esclusa is_offline_insert — non va sul server) am_cols = [r[1] for r in local_conn.execute("PRAGMA table_info(AnagraficaMedia)").fetchall() if r[1] != "is_offline_insert"] # Colonne LogAggiornamenti escluso id_aggiornamento (INTEGER PK auto-increment): # il server genera un nuovo ID autonomamente, evitando collisioni con record esistenti. log_cols = [r[1] for r in local_conn.execute("PRAGMA table_info(LogAggiornamenti)").fetchall() if r[1] != "id_aggiornamento"] server_conn = sqlite3.connect(str(server_path), timeout=30.0) try: for id_media in pending_ids: # Inserisci anagrafica — INSERT OR IGNORE: id_media è una chiave business # univoca (es. "FILM-1234567890"), nessun rischio di collisione con altri utenti. row = local_conn.execute( f"SELECT {', '.join(am_cols)} FROM AnagraficaMedia WHERE id_media = ?", (id_media,) ).fetchone() if row: placeholders = ", ".join(["?"] * len(am_cols)) server_conn.execute( f"INSERT OR IGNORE INTO AnagraficaMedia ({', '.join(am_cols)}) VALUES ({placeholders})", tuple(row) ) # Inserisci eventi — senza id_aggiornamento per evitare collisioni auto-increment. # Controllo idempotenza: salta se esiste già un evento con stessi campi chiave # (protezione in caso di retry dopo sync parziale). log_rows = local_conn.execute( f"SELECT {', '.join(log_cols)} FROM LogAggiornamenti WHERE id_media_fk = ?", (id_media,) ).fetchall() for log_row in log_rows: row_dict = dict(zip(log_cols, log_row)) # Verifica se l'evento esiste già sul server (stesso prodotto + tipo + data + contesto) exists = server_conn.execute(""" SELECT 1 FROM LogAggiornamenti WHERE id_media_fk = ? AND id_tipo_evento_fk = ? AND COALESCE(data_evento,'') = COALESCE(?,'') AND COALESCE(dettaglio_contesto,'') = COALESCE(?,'') LIMIT 1 """, ( row_dict.get("id_media_fk"), row_dict.get("id_tipo_evento_fk"), row_dict.get("data_evento"), row_dict.get("dettaglio_contesto"), )).fetchone() if not exists: placeholders = ", ".join(["?"] * len(log_cols)) server_conn.execute( f"INSERT INTO LogAggiornamenti ({', '.join(log_cols)}) VALUES ({placeholders})", tuple(log_row) ) server_conn.commit() logger.info(f"[sync] push_offline_inserts: {len(pending_ids)} prodotti sincronizzati sul server") finally: server_conn.close() # Clear flag is_offline_insert in locale placeholders = ", ".join(["?"] * len(pending_ids)) local_conn.execute( f"UPDATE AnagraficaMedia SET is_offline_insert = 0 WHERE id_media IN ({placeholders})", pending_ids ) local_conn.commit() # Cleanup: rimuovi backup offline e log JSONL (dati ora al sicuro sul server) _cleanup_offline_artifacts(local_db) local_conn.close() return len(pending_ids) except Exception as e: logger.error(f"[sync] push_offline_inserts fallito: {e}") return 0 def _cleanup_offline_artifacts(local_db: Path) -> None: """Rimuove backup offline e log JSONL dopo sync riuscito (dati al sicuro sul server).""" try: backup_dir = local_db.parent / "backups" if backup_dir.exists(): for f in backup_dir.glob("mediatrack_*.db"): f.unlink(missing_ok=True) log_file = local_db.parent / "offline_inserts_log.jsonl" if log_file.exists(): log_file.unlink(missing_ok=True) logger.info("[sync] Cleanup offline artifacts completato") except Exception as e: logger.warning(f"[sync] Cleanup offline artifacts fallito: {e}") def pull_ott_db() -> bool: """ Copia ott.db dal server in locale se il server ha una versione più recente. Flusso one-way: server → locale (nessun merge, nessun push). """ try: import json from app.config_loader import get_app_root config_path = get_app_root() / 'app' / 'modules' / 'linker' / 'config.json' if not config_path.exists(): return False cfg = json.loads(config_path.read_text(encoding='utf-8')) server_path = Path(cfg.get('ott_db_server_path', '')) local_rel = cfg.get('ott_db_local_path', 'local_db/ott.db') local_path = get_app_root() / local_rel if not server_path.exists(): logger.warning(f"[ott_sync] Server non raggiungibile: {server_path}") return False st = server_path.stat() if (local_path.exists() and local_path.stat().st_mtime >= st.st_mtime and local_path.stat().st_size == st.st_size): logger.debug("[ott_sync] ott.db già aggiornato — nessuna azione") return False import shutil local_path.parent.mkdir(parents=True, exist_ok=True) tmp = local_path.with_suffix('.tmp') shutil.copy2(str(server_path), str(tmp)) tmp.replace(local_path) logger.info("[ott_sync] Pull ott.db completato (server → locale)") return True except Exception as e: logger.error(f"[ott_sync] pull_ott_db fallito: {e}") return False def _background_push(): """Esegue push e reschedula il timer.""" push_to_server() _reschedule() def _reschedule(): """Reschedula il timer per il prossimo ciclo.""" global _timer try: sync_config, _ = _get_sync_config() interval = sync_config.get("interval_minutes", 5) * 60 except Exception: interval = 300 # fallback 5 minuti with _timer_lock: _timer = threading.Timer(interval, _background_push) _timer.daemon = True _timer.start() def start_background_sync(): """Avvia il timer di sync in background.""" logger.info("[sync] Background sync avviato") _reschedule() def stop_background_sync(): """Ferma il timer di sync.""" global _timer with _timer_lock: if _timer is not None: _timer.cancel() _timer = None logger.info("[sync] Background sync fermato")