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 :
${SLOW_SEARCH_URL}— catalogue principal (typiquementhttps://horizon.buchard.ch/api/travels).${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
idapparaît dans les deux sources, l'UPSERT de la deuxième source gagne etraw_data.isSeasideest positionné àtrue. C'est le comportement voulu — un voyage présent dans/seasidedoit ê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 :
| Chemin | Quand | Code |
|---|---|---|
| Démarrage du service | Au boot, si la dernière mise à jour > 15 min ou DB vide | BuchardDatabase._maybe_reindex() |
| Planificateur de fond | Toutes les REINDEX_FREQUENCY secondes (défaut 900s = 15min) | src/schedule_runner.start_schedule |
| Endpoint manuel | GET /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
_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/reindexne signale pas cet état au client — il répond toujours{"message": "Reindexing started", …}. - Le flag
_reindex_runningest remis àFalsedans lefinallyde_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)
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 += 1Garanties :
- 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ènename,subtitle,description,servicesIncluded,highLightsaprès nettoyage HTML (BeautifulSoup). Voirmodules/embeddings.md.- UPSERT idempotent : la PK
idest dérivée deraw_data->>'id', doncON CONFLICT (id)matche correctement. Les colonnes générées sont recalculées automatiquement.
Purge stale
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
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 travelsaprè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
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 retournenow()quandmax(updated_at)est NULL, doncdelta_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 quedb = BuchardDatabase()est au top-level). Les tests doivent injecterauto_reindex=False.
Échecs et observabilité
- Exceptions levées dans le thread : capturées par
_run_reindexet logguées vialogger.exception("Reindex failed"). Le flag_reindex_runningest 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é Horizon | Latence côté better-search |
|---|---|
| Création / édition d'un voyage | Jusqu'à REINDEX_FREQUENCY (défaut 15 min) |
| Publication d'une nouvelle occurrence | Idem |
| Suppression / dépublication | Idem (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é. Voirarchitecture/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.

