# ----------------------------------------------------------------------------- # Rent-Ka — Rental listings aggregator (Canada, outside Québec) # Author: 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, # cache des pages détail, cache de géocodage. # ----------------------------------------------------------------------------- from __future__ import annotations import json import sqlite3 import statistics import time from pathlib import Path from .schema import Listing DB_PATH = Path(__file__).resolve().parent.parent / "data" / "rentka.db" # Nombre d'exécutions consécutives où une annonce 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, unit_type TEXT, bedrooms REAL, bathrooms REAL, price REAL, price_label TEXT, availability TEXT, availability_date TEXT, area_sqft REAL, pets TEXT, furnished INTEGER, description TEXT, amenities TEXT, -- JSON (liste de textes source) details TEXT, -- JSON (champs structurés : inclusions, parking…) 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 ); 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_active ON listings(active); CREATE INDEX IF NOT EXISTS idx_listings_geo ON listings(lat, lng); 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é (NULL = retiré de l'affichage) ); CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid); 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 ); CREATE TABLE IF NOT EXISTS geocode_cache ( address TEXT PRIMARY KEY, -- adresse normalisée (clé de cache) lat REAL, lng REAL, provider TEXT, failed INTEGER DEFAULT 0, ts REAL ); CREATE TABLE IF NOT EXISTS users ( id INTEGER PRIMARY KEY AUTOINCREMENT, google_sub TEXT UNIQUE, -- identifiant stable Google (OpenID « sub ») ka_id TEXT, -- identifiant membre « ka-0123456789 » (unique) email TEXT, name TEXT, picture TEXT, -- URL de l'avatar Google public INTEGER DEFAULT 0, -- profil public /u/{ka_id} (opt-in) role TEXT, -- « locataire » | « gestionnaire » | NULL org_source TEXT, -- source réclamée (gestionnaire) created_at REAL, last_login REAL ); CREATE TABLE IF NOT EXISTS favorites ( user_id INTEGER NOT NULL, uid TEXT NOT NULL, -- listings.uid ts REAL, PRIMARY KEY (user_id, uid) ); CREATE TABLE IF NOT EXISTS image_checks ( url TEXT PRIMARY KEY, -- URL d'image vérifiée (imgcheck.py) status TEXT, -- ok | dead | small | error width INTEGER, height INTEGER, bytes INTEGER, sha1 TEXT, -- hash du contenu (doublons / placeholders) checked_at REAL ); CREATE TABLE IF NOT EXISTS env_tiles ( tile_key TEXT PRIMARY KEY, -- « ty,tx » (tuiles 0,5° — voir poi.py/environment.py) data TEXT, -- JSON : routes/rails/aéro/industriel/bars/cyclable/POI fetched_at REAL ); CREATE TABLE IF NOT EXISTS kascores ( coord_key TEXT PRIMARY KEY, -- « lat,lng » arrondi à 4 décimales (immeuble) lat REAL, lng REAL, walk REAL, -- KA Walk Score 0-100 (NULL = données insuffisantes) transit REAL, -- KA Transit Score (NULL = non desservi/inconnu) bike REAL, calme REAL, services REAL, global REAL, -- moyenne pondérée par défaut (voir kascores.py) details TEXT, -- JSON : détail par score (fiche + méthodologie) version TEXT, -- version du barème (recalcul si changement) computed_at REAL ); CREATE TABLE IF NOT EXISTS source_profiles ( source_id TEXT PRIMARY KEY, -- id de la source (data/sources.json) owner_user_id INTEGER, -- gestionnaire qui a réclamé la page tagline TEXT, -- accroche personnalisée description TEXT, website TEXT, phone TEXT, email TEXT, logo_url TEXT, updated_at REAL ); CREATE TABLE IF NOT EXISTS listing_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, uid TEXT NOT NULL, -- listings.uid ts REAL NOT NULL, event TEXT NOT NULL, -- description | superficie | dispo | -- inclusions | photos | disparition | -- reapparition (prix -> price_log) old TEXT, -- valeur avant (JSON compact) new TEXT -- valeur après (JSON compact) ); CREATE INDEX IF NOT EXISTS idx_levents_uid ON listing_events(uid, ts); CREATE TABLE IF NOT EXISTS buildings ( bkey TEXT PRIMARY KEY, -- clé d'adresse « civique|rue|ville » -- (dedup._parse_address) ou « geo:lat,lng » address TEXT, -- adresse d'affichage (la plus fréquente) city TEXT, lat REAL, lng REAL, stats TEXT, -- JSON : unités, annonces, loyers, -- rotation, pression, gestionnaires… version TEXT, -- version du calcul (building.py) computed_at REAL ); CREATE TABLE IF NOT EXISTS recycled_cache ( uid TEXT PRIMARY KEY, -- listings.uid matches TEXT, -- JSON : correspondances probables + signaux computed_at REAL ); CREATE TABLE IF NOT EXISTS managers ( id INTEGER PRIMARY KEY AUTOINCREMENT, source_id TEXT UNIQUE, -- id de la source (data/sources.json) name TEXT, -- nom canonique aliases TEXT, -- JSON : variantes de nom website TEXT, phone TEXT, -- résolution Google Maps (SerpApi) — voir rentka/managers.py gmaps_place_id TEXT, gmaps_data_id TEXT, gmaps_name TEXT, gmaps_address TEXT, gmaps_rating REAL, gmaps_reviews INTEGER, -- nombre d'avis annoncé par Google match_confidence REAL, -- 0..1 (association refusée si trop faible) match_signals TEXT, -- JSON : signaux ayant fondé l'association resolved_at REAL, resolve_failed INTEGER DEFAULT 0, review_stats TEXT, -- JSON : distribution, tendance, thèmes last_synced_at REAL ); CREATE TABLE IF NOT EXISTS manager_reviews ( id INTEGER PRIMARY KEY AUTOINCREMENT, manager_id INTEGER NOT NULL, external_review_id TEXT, -- id Google (déduplication) source TEXT DEFAULT 'google_maps', rating REAL, text TEXT, published_at TEXT, -- ISO si dérivable relative_date_raw TEXT, -- « il y a 3 mois » (texte source) author_name TEXT, author_review_count INTEGER, owner_response TEXT, owner_response_date TEXT, source_url TEXT, fetched_at REAL, content_hash TEXT, -- déduplication de secours analysis TEXT -- JSON : topics/sentiment/severity (managers.py) ); CREATE INDEX IF NOT EXISTS idx_mreviews_mgr ON manager_reviews(manager_id); CREATE UNIQUE INDEX IF NOT EXISTS idx_mreviews_ext ON manager_reviews(manager_id, external_review_id); """ # Colonnes ajoutées après la v1 — migration automatique des bases existantes. _MIGRATIONS = { "listings": { "availability_date": "TEXT", "area_sqft": "REAL", "pets": "TEXT", "furnished": "INTEGER", "details": "TEXT", "geocode_failed": "INTEGER DEFAULT 0", "miss_count": "INTEGER DEFAULT 0", "dauid": "TEXT", # aire de diffusion 2021 (stats de quartier) "digest": "TEXT", # JSON rentka/textmine.py (description structurée) "dup_of": "TEXT", # uid de l'annonce canonique si doublon inter-sources "dup_sources": "TEXT", # JSON : autres sources où l'annonce est publiée "bedrooms": "REAL", # chambres fermées (convention QC : 4½ = 2) "bathrooms": "REAL", # salles de bain (1.5 = salle d'eau en plus) "completeness": "REAL", # score de complétude 0–100 (quality.py) "published": "INTEGER DEFAULT 1", # 0 = quarantaine (sous le seuil) "quality_reasons": "TEXT", # JSON : raisons de la quarantaine "images_ok": "TEXT", # JSON : galerie nettoyée (imgcheck.py) "img_audit": "TEXT", # JSON : traçabilité du contrôle images "coord_key": "TEXT", # clé immeuble « lat,lng » 4 déc. (jointure kascores) "building_key": "TEXT", # clé immeuble par adresse (buildings.bkey) "province": "TEXT DEFAULT 'QC'", # « QC » | « ON » (expansion Ontario 2026-08) }, "sync_log": { "stats": "TEXT", }, "users": { "ka_id": "TEXT", "public": "INTEGER DEFAULT 0", "role": "TEXT", "org_source": "TEXT", "display_name": "TEXT", # nom choisi (prime sur le nom Google) "avatar_url": "TEXT", # photo téléversée (prime sur picture) "bio": "TEXT", "city": "TEXT", "phone": "TEXT", "website": "TEXT", "socials": "TEXT", # JSON {instagram, facebook, x, linkedin…} }, } def connect() -> sqlite3.Connection: DB_PATH.parent.mkdir(parents=True, exist_ok=True) con = sqlite3.connect(DB_PATH) con.row_factory = sqlite3.Row con.executescript(_SCHEMA) for table, cols in _MIGRATIONS.items(): existing = {r["name"] for r in con.execute(f"PRAGMA table_info({table})")} for col, decl in cols.items(): if col not in existing: con.execute(f"ALTER TABLE {table} ADD COLUMN {col} {decl}") # index sur des colonnes de migration : APRÈS les ALTER TABLE (une base # antérieure n'a pas encore la colonne au moment du CREATE INDEX) con.execute("CREATE UNIQUE INDEX IF NOT EXISTS users_ka_id ON users(ka_id)") con.execute("CREATE INDEX IF NOT EXISTS idx_listings_pub" " ON listings(active, published)") con.execute("CREATE INDEX IF NOT EXISTS idx_listings_coord" " ON listings(coord_key)") con.execute("CREATE INDEX IF NOT EXISTS idx_listings_bkey" " ON listings(building_key)") con.commit() # WAL : lectures (web) et écritures (sync, imgcheck) concurrentes sans # verrou global ; busy_timeout évite les « database is locked » ponctuels. con.execute("PRAGMA journal_mode=WAL") con.execute("PRAGMA busy_timeout=30000") 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). Retourne un message d'alerte, ou None si tout est normal. """ 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} annonce(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 annonces sans prix " f"(habituellement {statistics.median(rates):.0%}) — " "le format de la source a probablement changé") return None def _j(v) -> str: return json.dumps(v, ensure_ascii=False) def _log_events(con: sqlite3.Connection, uid: str, ts: float, old: sqlite3.Row, new: dict) -> None: """Journalise les changements observés d'une annonce (listing_events). Le prix a déjà son historique dédié (price_log) — ici : description, superficie, disponibilité, inclusions, photos. Valeurs compactes (pas de texte intégral) : l'objectif est la timeline de la fiche, pas l'archive. """ ev: list[tuple[str, str | None, str | None]] = [] old_desc = (old["description"] or "").strip() new_desc = (new.get("description") or "").strip() if old_desc != new_desc and (old_desc or new_desc): ev.append(("description", _j({"caracteres": len(old_desc)}), _j({"caracteres": len(new_desc)}))) if (old["area_sqft"] or None) != (new.get("area_sqft") or None): ev.append(("superficie", _j(old["area_sqft"]), _j(new.get("area_sqft")))) if (old["availability_date"] or None) != (new.get("availability_date") or None): ev.append(("dispo", _j(old["availability_date"]), _j(new.get("availability_date")))) try: inc_old = (json.loads(old["details"] or "{}") or {}).get("inclusions") or {} inc_new = (json.loads(new.get("details") or "{}") or {}).get("inclusions") or {} except (ValueError, TypeError): inc_old = inc_new = {} if inc_old != inc_new: keys = set(inc_old) | set(inc_new) diff_old = {k: inc_old.get(k) for k in keys if inc_old.get(k) != inc_new.get(k)} diff_new = {k: inc_new.get(k) for k in keys if inc_old.get(k) != inc_new.get(k)} if diff_old or diff_new: ev.append(("inclusions", _j(diff_old), _j(diff_new))) try: imgs_old = json.loads(old["images"] or "[]") imgs_new = json.loads(new.get("images") or "[]") except (ValueError, TypeError): imgs_old = imgs_new = [] if set(imgs_old) != set(imgs_new): ev.append(("photos", _j({"n": len(imgs_old)}), _j({"n": len(imgs_new)}))) if ev: con.executemany( "INSERT INTO listing_events (uid, ts, event, old, new)" " VALUES (?,?,?,?,?)", [(uid, ts, e, o, n) for e, o, n in ev]) def sync_source(con: sqlite3.Connection, source: str, listings: list[Listing]) -> dict: """Synchronise les annonces d'une source. - nouvelle annonce -> insertion - annonce modifiée -> mise à jour (comparaison de content_hash) - annonce disparue -> miss_count += 1, puis active=0 après MISS_GRACE exécutions consécutives (délai de grâce) - 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, images, description, area_sqft," " availability_date, details, active FROM listings WHERE uid=?", (lst.uid,)).fetchone() # score de complétude + seuil de publication (quarantaine sous le seuil) from .quality import evaluate q_score, q_pub, q_reasons = evaluate({ "title": lst.title, "price": lst.price, "address": lst.address, "city": lst.city, "sector": lst.sector, "lat": lst.lat, "lng": lst.lng, "description": lst.description, "images": lst.images, "unit_type": lst.unit_type, "bedrooms": lst.bedrooms, "bathrooms": lst.bathrooms, "availability": lst.availability, "availability_date": lst.availability_date, "area_sqft": lst.area_sqft, "amenities": lst.amenities, "details": lst.details, "url": lst.url, "province": getattr(lst, "province", None) or "QC"}) 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, province=(getattr(lst, "province", None) or "QC"), unit_type=lst.unit_type, bedrooms=lst.bedrooms, bathrooms=lst.bathrooms, price=lst.price, price_label=lst.price_label, availability=lst.availability, availability_date=lst.availability_date, area_sqft=lst.area_sqft, pets=lst.pets, furnished=(None if lst.furnished is None else int(lst.furnished)), description=lst.description, digest=(json.dumps(lst.digest, ensure_ascii=False) if getattr(lst, "digest", None) else None), amenities=json.dumps(lst.amenities, 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, completeness=q_score, published=int(q_pub), quality_reasons=(json.dumps(q_reasons, ensure_ascii=False) if q_reasons else None), ) if row is None: con.execute( """INSERT INTO listings (uid, source, external_id, url, title, address, sector, city, province, unit_type, bedrooms, bathrooms, price, price_label, availability, availability_date, area_sqft, pets, furnished, description, digest, amenities, details, images, lat, lng, content_hash, first_seen, last_seen, updated_at, miss_count, active, completeness, published, quality_reasons) VALUES (:uid,:source,:external_id,:url,:title,:address, :sector,:city,:province,:unit_type,:bedrooms,:bathrooms, :price,:price_label,:availability, :availability_date,:area_sqft,:pets,:furnished, :description,:digest,:amenities,:details,:images,:lat,:lng, :content_hash,:now,:now,:now,0,1, :completeness,:published,:quality_reasons)""", params) if lst.price is not None: # prix initial = point de départ de l'historique 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 des coordonnées géocodées par null con.execute( """UPDATE listings SET url=:url, title=:title, address=:address, sector=:sector, city=:city, province=:province, unit_type=:unit_type, bedrooms=:bedrooms, bathrooms=:bathrooms, price=:price, price_label=:price_label, availability=:availability, availability_date=:availability_date, area_sqft=:area_sqft, pets=:pets, furnished=:furnished, description=:description, digest=:digest, amenities=:amenities, 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, completeness=:completeness, published=:published, quality_reasons=:quality_reasons WHERE uid=:uid""", params) if (row["images"] or "[]") != params["images"]: # galerie modifiée : invalider l'audit d'images (imgcheck.py) con.execute("UPDATE listings SET images_ok=NULL, img_audit=NULL" " WHERE uid=?", (lst.uid,)) if lst.price != row["price"]: # changement de prix -> historique con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)", (lst.uid, now, lst.price)) # timeline de la fiche : changements observés (hors prix) _log_events(con, lst.uid, now, row, params) if not row["active"]: con.execute("INSERT INTO listing_events (uid, ts, event, old, new)" " VALUES (?,?,?,?,?)", (lst.uid, now, "reapparition", None, None)) updated += 1 else: if not row["active"]: con.execute("INSERT INTO listing_events (uid, ts, event, old, new)" " VALUES (?,?,?,?,?)", (lst.uid, now, "reapparition", None, None)) con.execute( "UPDATE listings SET last_seen=?, miss_count=0, active=1 WHERE uid=?", (now, lst.uid)) # Annonces de cette source qui n'apparaissent plus : délai de grâce, # puis désactivation. Suspendu si une dérive est 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"])) con.execute("INSERT INTO listing_events (uid, ts, event, old, new)" " VALUES (?,?,?,?,?)", (r["uid"], now, "disparition", None, None)) 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 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()