Skip to content

Architecture — Cycle de vie du reindex

Le reindex est le seul mécanisme qui écrit dans la table travels. Il pull l'amont (Horizon), calcule les embeddings, UPSERT, puis purge le stale. Cette page décrit son cycle complet et ses garanties.

Sources amont

Deux URLs sont pull, dans cet ordre :

  1. ${SLOW_SEARCH_URL} — catalogue principal (typiquement https://horizon.buchard.ch/api/travels).
  2. ${SLOW_SEARCH_URL}/seaside — voyages balnéaires.

Chaque URL est paginée (?size=10&page=1, puis page=2, etc.) jusqu'à recevoir une page avec totalFilteredRecords == 0.

L'ordre est important : /seaside est ingéré en dernier. Si un même id apparaît dans les deux sources, l'UPSERT de la deuxième source gagne et raw_data.isSeaside est positionné à true. C'est le comportement voulu — un voyage présent dans /seaside doit être marqué comme balnéaire.

Les records venant de /seaside se voient injecter record['isSeaside'] = True côté Python avant insertion, ce qui fait que la colonne générée is_seaside ressort à TRUE automatiquement (voir architecture/schema.md).

Déclenchement

Le reindex peut être lancé par trois chemins :

CheminQuandCode
Démarrage du serviceAu boot, si la dernière mise à jour > 15 min ou DB videBuchardDatabase._maybe_reindex()
Planificateur de fondToutes les REINDEX_FREQUENCY secondes (défaut 900s = 15min)src/schedule_runner.start_schedule
Endpoint manuelGET /reindex (curl, dashboard interne)server.reindex()db.reindex()

Tous ces chemins finissent dans la même fonction publique : src.reindex.reindex(pool).

Concurrence : lock global

python
_reindex_lock = Lock()
_reindex_running = False

def reindex(pool) -> bool:
    global _reindex_running
    with _reindex_lock:
        if _reindex_running:
            return False
        _reindex_running = True
    Thread(target=_run_reindex, args=(pool,), daemon=True).start()
    return True
  • Un seul reindex à la fois, peu importe le déclencheur.
  • Si un appel arrive pendant qu'un reindex tourne, l'appel retourne immédiatement (False) et le reindex en cours continue. L'endpoint /reindex ne signale pas cet état au client — il répond toujours {"message": "Reindexing started", …}.
  • Le flag _reindex_running est remis à False dans le finally de _run_reindex, même en cas d'exception.

Pipeline d'ingestion

┌──────────────────────────────────────────────────────────┐
│  _do_reindex(pool)                                       │
│                                                          │
│  time_start, new_ids = set(), base_url, seaside_url      │
│                                                          │
│  ─►  _ingest_source(pool, base_url, new_ids,             │
│                     mark_seaside=False)                  │
│                                                          │
│  ─►  _ingest_source(pool, seaside_url, new_ids,          │
│                     mark_seaside=True)                   │
│                                                          │
│  ─►  DELETE FROM travels WHERE id != ALL(new_ids)        │
│                                                          │
│  ─►  Métriques :                                         │
│        db_size.set(count)                                │
│        reindex_clock.set(now)                            │
│        reindex_timer.observe(elapsed)                    │
└──────────────────────────────────────────────────────────┘

_ingest_source(pool, url, new_ids, mark_seaside)

python
page = 1
while True:
    resp = requests.get(f"{url}?size=10&page={page}", timeout=_FETCH_TIMEOUT)
    resp.raise_for_status()
    result = resp.json()

    if result['totalFilteredRecords'] == 0:
        break    # source épuisée

    records = result['records']
    if mark_seaside:
        for record in records:
            record['isSeaside'] = True

    embedding_texts = [utils.get_embedding_text(r) for r in records]
    embeddings = infomaniak.get_embeddings(embedding_texts)  # batch

    with pool.connection() as conn:
        register_vector(conn)
        for record, embedding in zip(records, embeddings):
            new_ids.add(record['id'])
            conn.execute("""
                INSERT INTO travels (raw_data, embedding, updated_at)
                VALUES (%s, %s, NOW())
                ON CONFLICT (id) DO UPDATE SET
                    raw_data = EXCLUDED.raw_data,
                    embedding = EXCLUDED.embedding,
                    updated_at = NOW()
            """, (json.dumps(record), np.array(embedding, dtype=np.float32)))
        conn.commit()

    page += 1

Garanties :

  • Commit par page : si la page 3 plante, les pages 1 et 2 restent en DB. La purge stale ne tournera pas (cf. ci-dessous).
  • Batch d'embeddings = page entière (10 records max). Un seul appel Infomaniak par page.
  • get_embedding_text(record) concatène name, subtitle, description, servicesIncluded, highLights après nettoyage HTML (BeautifulSoup). Voir modules/embeddings.md.
  • UPSERT idempotent : la PK id est dérivée de raw_data->>'id', donc ON CONFLICT (id) matche correctement. Les colonnes générées sont recalculées automatiquement.

Purge stale

python
if new_ids:
    with pool.connection() as conn:
        conn.execute(
            "DELETE FROM travels WHERE id != ALL(%s)",
            (list(new_ids),)
        )
        conn.commit()

Note critique : la purge n'est exécutée que si les deux _ingest_source ont réussi. Si /seaside lève (timeout, 500, etc.), _do_reindex propage l'exception avant d'atteindre le DELETE. Les anciennes données restent intactes.

Ce comportement est testé dans tests/test_reindex.py::TestSeasideReindex::test_reindex_aborts_on_seaside_failure_keeps_old_data.

SLOW_SEARCH_URL avec ou sans trailing slash

python
base_url = os.getenv("SLOW_SEARCH_URL").rstrip("/")
seaside_url = f"{base_url}/seaside"

Le .rstrip("/") évite le double slash //seaside quand la variable d'env contient un trailing /. Testé.

Métriques émises

Voir modules/monitoring.md pour les détails :

  • reindex_time_s (histogram, buckets [15, 30, 45, 60, 75, 90, 120, 240]) — durée totale d'un reindex.
  • last_reindex_time (gauge) — timestamp Unix de la fin du dernier reindex.
  • db_size (gauge) — SELECT COUNT(*) FROM travels après le reindex.

Le last_reindex_time ne reflète pas la date côté DB. Pour ça, regarder MAX(updated_at) sur travels_view — c'est d'ailleurs ce qu'utilise _last_indexing_time() pour décider si un reindex au boot est nécessaire.

_maybe_reindex() au démarrage

python
def _maybe_reindex(self):
    delta_t = (datetime.now() - self._last_indexing_time()).total_seconds()
    if delta_t > 15 * 60 or delta_t < 1 or self.total_records == 0:
        self.reindex()
    else:
        logger.info("Last reindex was %s seconds ago. Skipping reindex", delta_t)
  • delta_t > 15 * 60 : données trop vieilles → reindex.
  • delta_t < 1 : DB vide (la fonction retourne now() quand max(updated_at) est NULL, donc delta_t ≈ 0) → reindex.
  • self.total_records == 0 : safety net redondant.
  • Sinon, on attend que le planificateur de fond lance le prochain cycle.

Cette logique tourne au moment de l'import du module src.db (parce que db = BuchardDatabase() est au top-level). Les tests doivent injecter auto_reindex=False.

Échecs et observabilité

  • Exceptions levées dans le thread : capturées par _run_reindex et logguées via logger.exception("Reindex failed"). Le flag _reindex_running est libéré, donc un reindex ultérieur peut redémarrer normalement.
  • Pas de retry automatique : un échec attend simplement le prochain cycle de 15 min (ou un appel manuel à /reindex).
  • Pas de circuit breaker sur Infomaniak : si l'API d'embedding est down, le reindex casse à la première page non vide. L'ancien contenu reste interrogeable, mais ne sera pas rafraîchi tant que l'amont n'est pas revenu.
  • Pas de transaction globale : commit par page. Si page 3 plante après commit de pages 1-2, la DB contient un mix (vieux records de pages 4+, nouveaux records de pages 1-2). La purge stale ne tournant pas, les anciens records hors-amont ne sont pas supprimés non plus. Au prochain reindex réussi, tout revient en ordre.

Côté front : quand voir des données fraîches

Action côté HorizonLatence côté better-search
Création / édition d'un voyageJusqu'à REINDEX_FREQUENCY (défaut 15 min)
Publication d'une nouvelle occurrenceIdem
Suppression / dépublicationIdem (le voyage disparaît au prochain DELETE)
Modif marketing (description, photos)Idem (texte embeddé → vecteur recalculé)

Pour un rafraîchissement immédiat : curl https://better-search-host/reindex (réponse 200 immédiate, reindex en arrière-plan).

Ce qui n'est PAS fait

  • Pas de diff incrémental. Chaque reindex re-télécharge tout le catalogue (~quelques centaines de records) et recalcule tous les embeddings. À l'échelle actuelle (< 1000 voyages, < 1 min de cycle), ce n'est pas un problème.
  • Pas de versioning des embeddings. Si le modèle BGE change ou si on bascule sur un autre fournisseur, il faut une migration qui change vector(3584) → nouvelle dimension et un reindex forcé. Voir architecture/adr/0003-bge-multilingual-gemma2.md.
  • Pas de hash content-based pour éviter de re-embedder un record inchangé. Trivial à ajouter (sha256(raw_data) stocké en colonne, skip si égal), pas implémenté car le coût total d'embedding sur le catalogue est négligeable.

Contributors

No contributors

Changelog

No recent changes