# Author: Simon-Pierre Boucher # Contact: contact@spboucher.ai # Project: Toit-Ka # ----------------------------------------------------------------------------- # replicate.py : réplication LECTURE SEULE des BD de prod vers ce nœud. # # Lou-Ka tourne sur M3U96b (~/apps/lou-ka/data/louka.db) # Immo-Ka tourne sur M4M64a (~/apps/immo-ka/data/immoka.db) # Toit-Ka tourne sur M3U96a et TIRE des snapshots cohérents via le LAN : # 1. ssh "sqlite3 '.backup /tmp/-snap.db'" (snapshot sûr # même en cours d'écriture — jamais un rsync direct d'une DB WAL vivante) # 2. rsync du snapshot vers data/replicas/.db (échange atomique) # # Config .env (absente en dev -> replicate ne fait rien, l'ETL lit les BD # locales pointées par TOITKA_LOUKA_DB / TOITKA_IMMOKA_DB) : # TOITKA_REPLICA_SOURCES=louka=192.168.2.82:apps/lou-ka/data/louka.db,immoka=192.168.2.83:apps/immo-ka/data/immoka.db # TOITKA_REPLICA_USER=simon-pierreboucher # ----------------------------------------------------------------------------- from __future__ import annotations import os import subprocess import time from pathlib import Path from .db import DATA REPLICAS = DATA / "replicas" SSH_OPTS = ["-o", "BatchMode=yes", "-o", "ConnectTimeout=10", "-o", "StrictHostKeyChecking=accept-new"] def _sources() -> list[tuple[str, str, str]]: """[(nom, hôte, chemin distant relatif au $HOME), …] depuis l'environnement.""" raw = os.environ.get("TOITKA_REPLICA_SOURCES", "").strip() out = [] for part in filter(None, (p.strip() for p in raw.split(","))): name, _, rest = part.partition("=") host, _, path = rest.partition(":") if name and host and path: out.append((name.strip(), host.strip(), path.strip())) return out def _pull(name: str, host: str, remote_db: str, user: str) -> bool: t0 = time.time() snap = f"/tmp/toitka-{name}-snap.db" dest = REPLICAS / f"{name}.db" tmp = REPLICAS / f"{name}.db.part" # 1. snapshot cohérent côté source (sqlite .backup gère le WAL en cours) r = subprocess.run( ["ssh", *SSH_OPTS, f"{user}@{host}", f"sqlite3 ~/{remote_db} \".backup {snap}\""], capture_output=True, text=True, timeout=600) if r.returncode != 0: print(f"[replicate] {name} : snapshot ÉCHOUÉ — {r.stderr.strip()[:300]}") return False # 2. transfert (rsync incrémental) puis échange atomique r = subprocess.run( ["rsync", "-a", "--partial", "-e", "ssh " + " ".join(SSH_OPTS), f"{user}@{host}:{snap}", str(tmp)], capture_output=True, text=True, timeout=1800) if r.returncode != 0 or not tmp.exists(): print(f"[replicate] {name} : rsync ÉCHOUÉ — {r.stderr.strip()[:300]}") tmp.unlink(missing_ok=True) return False os.replace(tmp, dest) subprocess.run(["ssh", *SSH_OPTS, f"{user}@{host}", f"rm -f {snap}"], capture_output=True, timeout=60) mb = dest.stat().st_size / 1e6 print(f"[replicate] {name} : {mb:.0f} Mo en {time.time() - t0:.1f} s") return True def run() -> None: sources = _sources() if not sources: return # dev : rien à répliquer, l'ETL lit en local REPLICAS.mkdir(parents=True, exist_ok=True) user = os.environ.get("TOITKA_REPLICA_USER", os.environ.get("USER", "")) for name, host, path in sources: try: _pull(name, host, path, user) except subprocess.TimeoutExpired: print(f"[replicate] {name} : délai dépassé") except Exception as e: print(f"[replicate] {name} : {type(e).__name__}: {e}")