# ----------------------------------------------------------------------------- # Immo-Ka — Agrégateur de maisons à vendre (province de Québec) # Auteur : Simon-Pierre Boucher — contact@spboucher.ai # db.py : persistance SQLite — upsert avec détection de changements, # cycle de vie avec délai de grâce (2 syncs), détection de dérive, # historique de prix, cache des pages détail. # ----------------------------------------------------------------------------- from __future__ import annotations import json import sqlite3 import statistics import time from pathlib import Path from .schema import PropertyListing DB_PATH = Path(__file__).resolve().parent.parent / "data" / "immoka.db" # Nombre d'exécutions consécutives où une propriété doit être absente de la # source avant d'être désactivée (délai de grâce contre les ratés ponctuels). MISS_GRACE = 2 # Dérive : si une source retourne <= DRIFT_RATIO × sa médiane historique # (médiane >= DRIFT_MIN_BASE annonces), on alerte et on suspend les retraits. DRIFT_RATIO = 0.25 DRIFT_MIN_BASE = 8 DRIFT_HISTORY = 5 _SCHEMA = """ CREATE TABLE IF NOT EXISTS listings ( uid TEXT PRIMARY KEY, source TEXT NOT NULL, external_id TEXT NOT NULL, url TEXT, title TEXT, address TEXT, sector TEXT, city TEXT, region TEXT, property_type TEXT, price REAL, price_label TEXT, bedrooms INTEGER, bathrooms INTEGER, powder_rooms INTEGER, area_sqft REAL, lot_sqft REAL, year_built INTEGER, mls TEXT, status TEXT DEFAULT 'a-vendre', broker_name TEXT, broker_phone TEXT, agency TEXT, description TEXT, features TEXT, -- JSON (liste de textes source) details TEXT, -- JSON (champs structurés) images TEXT, -- JSON lat REAL, lng REAL, geocode_failed INTEGER DEFAULT 0, content_hash TEXT, first_seen REAL, last_seen REAL, updated_at REAL, miss_count INTEGER DEFAULT 0, active INTEGER DEFAULT 1, dup_hidden INTEGER DEFAULT 0 -- 1 = doublon de sous-agence masqué (dédup Centris) ); CREATE INDEX IF NOT EXISTS idx_listings_source ON listings(source); CREATE INDEX IF NOT EXISTS idx_listings_city ON listings(city); CREATE INDEX IF NOT EXISTS idx_listings_type ON listings(property_type); CREATE INDEX IF NOT EXISTS idx_listings_active ON listings(active); CREATE INDEX IF NOT EXISTS idx_listings_extid ON listings(external_id); CREATE TABLE IF NOT EXISTS sync_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, source TEXT, ts REAL, found INTEGER, added INTEGER, updated INTEGER, removed INTEGER, ok INTEGER, message TEXT, stats TEXT -- JSON : taux de champs null, missed, alerte… ); CREATE TABLE IF NOT EXISTS detail_cache ( source TEXT NOT NULL, external_id TEXT NOT NULL, key TEXT, -- hash du contenu « liste » de l'annonce payload TEXT, -- JSON opaque propre au connecteur fetched_at REAL, PRIMARY KEY (source, external_id) ); CREATE TABLE IF NOT EXISTS price_log ( uid TEXT NOT NULL, ts REAL NOT NULL, price REAL -- prix observé (baisses/hausses de prix demandé) ); CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid); CREATE TABLE IF NOT EXISTS geocode_cache ( address TEXT PRIMARY KEY, lat REAL, lng REAL, provider TEXT, failed INTEGER DEFAULT 0, ts REAL, muni TEXT -- municipalité officielle (clés « ville:… ») ); CREATE TABLE IF NOT EXISTS poi_cache ( coord_key TEXT PRIMARY KEY, -- "lat,lng" arrondi à 4 décimales (~11 m) lat REAL, lng REAL, pois TEXT, -- JSON : [{cat, name, dist_m}] (plus proche/catégorie) ts REAL ); """ _SCHEMA_READY = False # le schéma/migration ne s'exécute qu'UNE fois par process def _init_schema(con: sqlite3.Connection) -> None: """Création de schéma + migration — coûteuse (write-lock). À faire une seule fois par process : l'exécuter à chaque connexion sérialisait les requêtes du web sur un verrou d'écriture (deadlock/starvation sous charge).""" con.executescript(_SCHEMA) cols = {r["name"] for r in con.execute("PRAGMA table_info(listings)")} if "agency" not in cols: con.execute("ALTER TABLE listings ADD COLUMN agency TEXT") if "dup_hidden" not in cols: con.execute("ALTER TABLE listings ADD COLUMN dup_hidden INTEGER DEFAULT 0") if "dup_of" not in cols: # uid de la fiche visible au profit de laquelle con.execute("ALTER TABLE listings ADD COLUMN dup_of TEXT") # celle-ci est masquée if "dauid" not in cols: # aire de diffusion 2021 (stats de quartier) con.execute("ALTER TABLE listings ADD COLUMN dauid TEXT") if "vraiprix" not in cols: # estimation Vrai-Prix (JSON) + lien analyse con.execute("ALTER TABLE listings ADD COLUMN vraiprix TEXT") gcols = {r["name"] for r in con.execute("PRAGMA table_info(geocode_cache)")} if "muni" not in gcols: # municipalité officielle (entrées « ville:… ») con.execute("ALTER TABLE geocode_cache ADD COLUMN muni TEXT") con.execute("CREATE INDEX IF NOT EXISTS idx_listings_duphidden ON listings(dup_hidden)") con.execute("CREATE INDEX IF NOT EXISTS idx_listings_dupof ON listings(dup_of)") con.execute("CREATE INDEX IF NOT EXISTS idx_listings_geo ON listings(lat, lng)") con.commit() def refresh_dedup(con: sqlite3.Connection) -> int: """Pré-calcule la déduplication de la famille sous-agences (`*_ag_*`) dans la colonne `dup_hidden`, pour que la lecture soit instantanée (« AND dup_hidden=0 ») au lieu d'un sous-select corrélé par ligne (300 s sur 75 k lignes). Règle (identique à l'ancienne clause) : une fiche de sous-agence est masquée si une fiche de plus haute priorité — même n° Centris (external_id), active — existe : le flux central/agence principale d'abord, sinon la sous-agence au plus petit uid. Les fiches non sous-agence ne sont jamais masquées par cette règle. Une 3e passe (dedup_by_address) masque ensuite les doublons INTER-SOURCES sans n° Centris commun (même adresse + type + prix ±1 %). Retourne le nombre total de fiches masquées.""" AG = "source LIKE '%\\_ag\\_%' ESCAPE '\\'" NOTAG = "source NOT LIKE '%\\_ag\\_%' ESCAPE '\\'" con.execute("UPDATE listings SET dup_hidden=0, dup_of=NULL") # 1) masquer les sous-agences dont le n° Centris est porté par une fiche # canonique (non sous-agence) active — semi-jointure, rapide. # dup_of = la fiche canonique (pour la box « Aussi publiée sur… »). con.execute( f"UPDATE listings SET dup_hidden=1," f" dup_of=(SELECT MIN(d.uid) FROM listings d WHERE d.active=1" f" AND d.external_id=listings.external_id AND d.{NOTAG})" f" WHERE active=1 AND {AG}" f" AND external_id IN (SELECT external_id FROM listings" f" WHERE active=1 AND {NOTAG})") # 2) parmi les sous-agences restantes (sans canonique), ne garder que le plus # petit uid par n° Centris. con.execute( f"UPDATE listings SET dup_hidden=1," f" dup_of=(SELECT MIN(d.uid) FROM listings d WHERE d.active=1" f" AND d.external_id=listings.external_id AND d.{AG} AND d.dup_hidden=0)" f" WHERE active=1 AND {AG} AND dup_hidden=0" f" AND uid > (SELECT MIN(d.uid) FROM listings d WHERE d.active=1" f" AND d.external_id=listings.external_id AND d.{AG} AND d.dup_hidden=0)") # 3) dédup INTER-SOURCES par adresse : la même propriété publiée sur deux # plateformes (ex. courtier + Kijiji) sans n° Centris commun. Règle # conservatrice : même adresse normalisée (civique + rue + ville) + même # type + prix identique à ±1 % → on garde la fiche la plus autoritaire. try: n_addr = dedup_by_address(con) if n_addr: print(f"[immo-ka] dédup adresse: {n_addr} doublon(s) inter-sources masqué(s)") except Exception: import traceback traceback.print_exc() n = con.execute("SELECT COUNT(*) c FROM listings WHERE dup_hidden=1").fetchone()["c"] con.commit() return n # petites annonces généralistes (republication d'annonces d'ailleurs) : moins # autoritaires que la source primaire (courtier / FSBO première main). # fb_marketplace : fiches anonymes/republication — jamais préférées à un courtier. _PETITES_ANNONCES = {"kijiji", "lespac", "fb_marketplace"} def _source_rank(source: str) -> int: """Autorité d'une source pour la dédup d'adresse : 0 = source primaire (bannière/agence/FSBO), 1 = sous-agence (_ag_), 2 = petites annonces (republication probable).""" if source in _PETITES_ANNONCES: return 2 if "_ag_" in source: return 1 return 0 def dedup_by_address(con: sqlite3.Connection) -> int: """Passe de déduplication conservatrice par ADRESSE (inter-sources). Clé : n° civique(s) + mots significatifs de la rue + ville normalisée + type canonique + n° d'app/unité (vide s'il n'y en a pas — deux unités différentes d'un même immeuble ne partagent JAMAIS la même clé, et une adresse sans unité ne s'apparie pas à une adresse avec unité). Dans une clé, seules les fiches au prix identique à ±1 % sont considérées comme doublons ; si une même source apparaît deux fois dans le groupe (probables unités jumelles d'un projet neuf), le groupe ENTIER est ignoré. On garde la fiche la plus autoritaire (source primaire > sous-agence _ag_ > petites annonces), puis à autorité égale celle qui a un COURTIER/agence (jamais une fiche anonyme devant un courtier), puis le plus petit uid. Retourne le nb masqué.""" import re as _re from .vraiprix_local import _addr_parts, _norm, _muni_norm, _APP_RE groups: dict[tuple, list] = {} for r in con.execute( "SELECT uid, source, address, city, price, property_type," " broker_name, agency" " FROM listings WHERE active=1 AND dup_hidden=0 AND address<>''" " AND city<>'' AND price IS NOT NULL AND property_type<>''"): civs, words = _addr_parts(r["address"]) if not civs or not words: continue a = _norm(r["address"]).split(",")[0] mapt = _re.search(_APP_RE, a) apt = "" if mapt: toks = _re.findall(r"[a-z0-9]+", mapt.group(0)) # dernier token = le n° d'unité (le 1er est le mot-clé app/unité/#) apt = toks[-1] if toks else "" key = ("-".join(civs), " ".join(sorted(set(words))), _muni_norm(r["city"]), _norm(r["property_type"]), apt) anonyme = 0 if (r["broker_name"] or r["agency"]) else 1 groups.setdefault(key, []).append( (r["price"], _source_rank(r["source"]), anonyme, r["uid"], r["source"])) hidden = 0 for rows in groups.values(): if len(rows) < 2: continue rows.sort() # par prix croissant cluster: list = [] for row in rows: if cluster and row[0] > cluster[0][0] * 1.01: hidden += _mask_cluster(con, cluster) cluster = [] cluster.append(row) hidden += _mask_cluster(con, cluster) return hidden def _mask_cluster(con: sqlite3.Connection, cluster: list) -> int: """Masque les doublons d'un groupe (même clé d'adresse, prix ±1 %), puis FUSIONNE en « golden record » : les champs vides de la fiche conservée sont complétés depuis les doublons masqués (description plus longue, galerie plus riche, superficies, année, GPS, téléphone) — le meilleur des deux sources.""" if len(cluster) < 2: return 0 sources = [c[4] for c in cluster] if len(set(sources)) != len(sources): return 0 # même source en double = probables unités distinctes : prudence keep = min(cluster, key=lambda c: (c[1], c[2], c[3])) # (autorité, anonyme, uid) n = 0 donors = [] for c in cluster: if c[3] != keep[3]: con.execute("UPDATE listings SET dup_hidden=1, dup_of=? WHERE uid=?", (keep[3], c[3])) donors.append(c[3]) n += 1 try: _merge_golden(con, keep[3], donors) except Exception: pass # la fusion est un bonus — ne jamais casser la dédup return n _MERGE_NUM_FIELDS = ("bedrooms", "bathrooms", "powder_rooms", "area_sqft", "lot_sqft", "year_built") def _merge_golden(con: sqlite3.Connection, keep_uid: str, donor_uids: list[str]) -> None: """Complète les champs vides de `keep_uid` depuis ses doublons masqués.""" if not donor_uids: return cols = ("uid, description, images, bedrooms, bathrooms, powder_rooms," " area_sqft, lot_sqft, year_built, lat, lng, broker_phone") keep = con.execute(f"SELECT {cols} FROM listings WHERE uid=?", (keep_uid,)).fetchone() if keep is None: return sets, args = [], [] kimgs = len(json.loads(keep["images"] or "[]")) kdesc = len(keep["description"] or "") best_imgs, best_desc = None, None donor_rows = con.execute( f"SELECT {cols} FROM listings WHERE uid IN " f"({','.join('?' * len(donor_uids))})", donor_uids).fetchall() merged_num: dict = {} for d in donor_rows: di = json.loads(d["images"] or "[]") if len(di) > max(kimgs, len(json.loads(best_imgs or "[]"))): best_imgs = d["images"] dd = d["description"] or "" if len(dd) > max(kdesc, 80, len(best_desc or "")): best_desc = dd for f in _MERGE_NUM_FIELDS: if keep[f] is None and merged_num.get(f) is None and d[f] is not None: merged_num[f] = d[f] # GPS : toujours la PAIRE du même donneur (jamais lat et lng mélangés) if (keep["lat"] is None and "lat" not in merged_num and d["lat"] is not None and d["lng"] is not None): merged_num["lat"], merged_num["lng"] = d["lat"], d["lng"] if not keep["broker_phone"] and d["broker_phone"] and "broker_phone" not in merged_num: merged_num["broker_phone"] = d["broker_phone"] if best_imgs is not None: sets.append("images=?"); args.append(best_imgs) if best_desc is not None and kdesc < 80: sets.append("description=?"); args.append(best_desc) for f, v in merged_num.items(): sets.append(f"{f}=?"); args.append(v) if sets: args.append(keep_uid) con.execute(f"UPDATE listings SET {', '.join(sets)} WHERE uid=?", args) def connect() -> sqlite3.Connection: global _SCHEMA_READY DB_PATH.parent.mkdir(parents=True, exist_ok=True) con = sqlite3.connect(DB_PATH, timeout=60) con.row_factory = sqlite3.Row # WAL + busy_timeout : accès concurrents (watcher + web) sans « db is locked ». con.execute("PRAGMA journal_mode=WAL") # 120 s : couvre les longues transactions (refresh_dedup ~70 s) sans que les # autres écrivains (géocodage, vraiprix) ne lèvent « database is locked ». con.execute("PRAGMA busy_timeout=120000") con.execute("PRAGMA synchronous=NORMAL") if not _SCHEMA_READY: _init_schema(con) _SCHEMA_READY = True return con # --------------------------------------------------------------------------- # Synchronisation d'une source # --------------------------------------------------------------------------- def _drift_alert(con: sqlite3.Connection, source: str, found: int, null_price_rate: float) -> str | None: """Détecte une dérive du connecteur (chute du volume ou des prix extraits).""" hist = con.execute( "SELECT found, stats FROM sync_log WHERE source=? AND ok=1" " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall() if len(hist) < 3: return None med_found = statistics.median(r["found"] for r in hist) if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found: return (f"dérive: {found} propriété(s) trouvée(s) contre une médiane de " f"{med_found:.0f} — retraits suspendus, vérifier le connecteur") if found >= DRIFT_MIN_BASE and null_price_rate >= 0.8: rates = [] for r in hist: try: rates.append(json.loads(r["stats"] or "{}")["null_price_rate"]) except (KeyError, ValueError, TypeError): continue if rates and statistics.median(rates) <= 0.3: return (f"dérive: {null_price_rate:.0%} des propriétés sans prix " f"(habituellement {statistics.median(rates):.0%}) — " "le format de la source a probablement changé") return None def sync_source(con: sqlite3.Connection, source: str, listings: list[PropertyListing]) -> dict: """Synchronise les propriétés d'une source. - nouvelle propriété -> insertion - propriété modifiée -> mise à jour (comparaison de content_hash) - propriété disparue -> miss_count += 1, puis active=0 après MISS_GRACE exécutions consécutives (vendue ou retirée) - dérive détectée -> alerte consignée, retraits suspendus """ now = time.time() added = updated = 0 seen_uids = set() n = len(listings) null_price = sum(1 for l in listings if l.price is None) null_addr = sum(1 for l in listings if not l.address) null_price_rate = round(null_price / n, 3) if n else 0.0 alert = _drift_alert(con, source, n, null_price_rate) for lst in listings: seen_uids.add(lst.uid) h = lst.content_hash() row = con.execute("SELECT content_hash, price FROM listings WHERE uid=?", (lst.uid,)).fetchone() params = dict( uid=lst.uid, source=lst.source, external_id=lst.external_id, url=lst.url, title=lst.title, address=lst.address, sector=lst.sector, city=lst.city, region=lst.region, property_type=lst.property_type, price=lst.price, price_label=lst.price_label, bedrooms=lst.bedrooms, bathrooms=lst.bathrooms, powder_rooms=lst.powder_rooms, area_sqft=lst.area_sqft, lot_sqft=lst.lot_sqft, year_built=lst.year_built, mls=lst.mls, status=lst.status, broker_name=lst.broker_name, broker_phone=lst.broker_phone, agency=lst.agency, description=lst.description, features=json.dumps(lst.features, ensure_ascii=False), details=json.dumps(lst.details, ensure_ascii=False), images=json.dumps(lst.images, ensure_ascii=False), lat=lst.lat, lng=lst.lng, content_hash=h, now=now, ) if row is None: con.execute( """INSERT INTO listings (uid, source, external_id, url, title, address, sector, city, region, property_type, price, price_label, bedrooms, bathrooms, powder_rooms, area_sqft, lot_sqft, year_built, mls, status, broker_name, broker_phone, agency, description, features, details, images, lat, lng, content_hash, first_seen, last_seen, updated_at, miss_count, active) VALUES (:uid,:source,:external_id,:url,:title,:address, :sector,:city,:region,:property_type,:price,:price_label, :bedrooms,:bathrooms,:powder_rooms,:area_sqft,:lot_sqft, :year_built,:mls,:status,:broker_name,:broker_phone, :agency,:description,:features,:details,:images,:lat,:lng, :content_hash,:now,:now,:now,0,1)""", params) if lst.price is not None: con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)", (lst.uid, now, lst.price)) added += 1 elif row["content_hash"] != h: # COALESCE : ne jamais écraser par null des coordonnées géocodées ni # les année/superficies remplies par enrichissement (rôle d'évaluation, # fiche détail) quand la source liste ne les fournit pas con.execute( """UPDATE listings SET url=:url, title=:title, address=:address, sector=:sector, city=:city, region=:region, property_type=:property_type, price=:price, price_label=:price_label, bedrooms=:bedrooms, bathrooms=:bathrooms, powder_rooms=:powder_rooms, area_sqft=COALESCE(:area_sqft, area_sqft), lot_sqft=COALESCE(:lot_sqft, lot_sqft), year_built=COALESCE(:year_built, year_built), mls=:mls, status=:status, broker_name=:broker_name, broker_phone=:broker_phone, agency=:agency, description=:description, features=:features, details=:details, images=:images, lat=COALESCE(:lat, lat), lng=COALESCE(:lng, lng), content_hash=:content_hash, last_seen=:now, updated_at=:now, miss_count=0, active=1 WHERE uid=:uid""", params) if lst.price != row["price"]: # baisse/hausse de prix -> historique con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)", (lst.uid, now, lst.price)) updated += 1 else: con.execute( "UPDATE listings SET last_seen=?, miss_count=0, active=1 WHERE uid=?", (now, lst.uid)) # Propriétés de cette source qui n'apparaissent plus : délai de grâce, # puis désactivation (vendue/retirée). Suspendu si dérive détectée. removed = missed = 0 if not alert: for r in con.execute( "SELECT uid, miss_count FROM listings WHERE source=? AND active=1", (source,)).fetchall(): if r["uid"] in seen_uids: continue missed += 1 if r["miss_count"] + 1 >= MISS_GRACE: con.execute( "UPDATE listings SET active=0, miss_count=?, updated_at=?" " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"])) removed += 1 else: con.execute("UPDATE listings SET miss_count=miss_count+1 WHERE uid=?", (r["uid"],)) stats = { "null_price_rate": null_price_rate, "null_address_rate": round(null_addr / n, 3) if n else 0.0, "missed": missed, } if alert: stats["alert"] = alert con.execute( "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok," " message, stats) VALUES (?,?,?,?,?,?,1,?,?)", (source, now, n, added, updated, removed, alert or "ok", json.dumps(stats, ensure_ascii=False))) con.commit() out = {"source": source, "found": n, "added": added, "updated": updated, "removed": removed} if alert: out["alert"] = alert return out def log_failure(con: sqlite3.Connection, source: str, message: str) -> None: con.execute( "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)" " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message)) con.commit() # --------------------------------------------------------------------------- # Cache des pages détail (« détail si nouveau/modifié ») # --------------------------------------------------------------------------- def get_cached_detail(con: sqlite3.Connection, source: str, external_id: str, key: str) -> dict | None: """Payload détail mis en cache si la clé (hash liste) n'a pas changé.""" row = con.execute( "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?", (source, external_id)).fetchone() if row and row["key"] == key and row["payload"]: try: return json.loads(row["payload"]) except ValueError: return None return None def get_stale_detail(con: sqlite3.Connection, source: str, external_id: str) -> dict | None: """Payload détail SANS vérifier la clé — repli « périmé plutôt que rien » quand le budget de re-fetch d'un cycle est épuisé (ex. bump de version de clé) : la fiche garde photos/détails existants en attendant son re-parse.""" row = con.execute( "SELECT payload FROM detail_cache WHERE source=? AND external_id=?", (source, external_id)).fetchone() if row and row["payload"]: try: return json.loads(row["payload"]) except ValueError: return None return None def put_cached_detail(con: sqlite3.Connection, source: str, external_id: str, key: str, payload: dict) -> None: con.execute( "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)" " VALUES (?,?,?,?,?)" " ON CONFLICT(source, external_id) DO UPDATE SET" " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at", (source, external_id, key, json.dumps(payload, ensure_ascii=False), time.time())) con.commit()