Module — Job de reindex
Le job qui pull l'amont Horizon, embed, UPSERT, purge stale. Code dans src/reindex.py (130 lignes) et src/schedule_runner.py (32 lignes). Cycle complet documenté en architecture : architecture/reindex-lifecycle.md. Cette page se concentre sur le point de vue module : APIs publiques, points d'extension, debug.
Surface publique
# src.reindex
def reindex(pool) -> bool:
"""Start a background reindex. Returns False if one is already running."""
def _do_reindex(pool):
"""Synchronous reindex logic. Can be called directly in tests."""# src.schedule_runner
def start_schedule():
"""Lance un thread qui boucle sur reindex() toutes REINDEX_FREQUENCY secondes."""# src.db
class BuchardDatabase:
def reindex(self):
"""Threaded reindex job — proxy vers reindex.reindex(self._pool)."""# src.server
@app.get("/reindex")
def reindex():
db.reindex()
return {"message": "Reindexing started", "time": …}Point d'entrée
┌──────────────────────────────┐
│ app.py │
│ ─ assert env vars │
│ ─ load DATABASE_URL_FILE │
│ ─ start_schedule() │
└──────────────┬───────────────┘
│
▼
┌──────────────────────────────┐
│ schedule_runner.start_… │
│ Thread loop infinie: │
│ ─ reindex() │
│ ─ sleep REINDEX_FREQUENCY │
└──────────────┬───────────────┘
│
▼
┌──────────────────────────────┐
│ src.server.reindex() │
│ → db.reindex() │
│ → reindex_mod.reindex(pool) │
│ → Thread(_run_reindex) │
│ → _do_reindex(pool) │
└──────────────────────────────┘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
def _run_reindex(pool):
global _reindex_running
try:
_do_reindex(pool)
except Exception:
logger.exception("Reindex failed")
finally:
with _reindex_lock:
_reindex_running = False_reindex_runningest lu/écrit sous le lock pour éviter une race entrewith _reindex_locketThread(...).start().daemon=True: le thread meurt avec le process. Pas de orphan thread au shutdown.except Exception: capture toutes les erreurs pour ne pas laisser le flag àTruead vitam.
Pipeline d'ingestion (_do_reindex)
Cf. architecture/reindex-lifecycle.md pour le détail. Résumé :
def _do_reindex(pool):
time_start = time.time()
new_ids = set()
base_url = os.getenv("SLOW_SEARCH_URL").rstrip("/")
seaside_url = f"{base_url}/seaside"
_ingest_source(pool, base_url, new_ids, mark_seaside=False)
_ingest_source(pool, seaside_url, new_ids, mark_seaside=True)
if new_ids:
with pool.connection() as conn:
conn.execute("DELETE FROM travels WHERE id != ALL(%s)", (list(new_ids),))
conn.commit()
# Métriques
db_size.set(count)
reindex_clock.set(now)
reindex_timer.observe(elapsed)_ingest_source(pool, url, new_ids, mark_seaside)
- Pagine
url?size=10&page=1..njusqu'àtotalFilteredRecords == 0. - Pour chaque page :
- Si
mark_seaside, injecterecord['isSeaside'] = True(pour la colonne généréeis_seaside). - Calcule les embeddings via
infomaniak.get_embeddings(texts)(batch). INSERT … ON CONFLICT (id) DO UPDATEsur chaque record.- Commit par page.
- Si
- Exception propagée — le caller (
_do_reindex) ne fait PAS le DELETE final.
Garanties
- Idempotent : un reindex sur une DB déjà à jour est un no-op fonctionnel (mêmes UPSERT, embedding recalculé).
- Atomic par page : commit par page, donc une page intermédiaire qui plante laisse les précédentes en DB.
- Pas de purge si exception :
DELETEn'est atteint que si les deux sources se sont bien terminées.
Configuration
| Env var | Défaut | Effet |
|---|---|---|
REINDEX_FREQUENCY | 900 | Secondes entre cycles auto (entier strict) |
SLOW_SEARCH_URL | obligatoire | Base URL de l'amont, sans trailing / requis |
Constantes en code :
| Constante | Valeur | Effet |
|---|---|---|
_FETCH_TIMEOUT | 30 s | Timeout HTTP requests.get |
| Page size amont | 10 | Hardcodé dans le ?size=10 de l'URL |
Métriques émises
(Voir modules/monitoring.md pour le détail Prometheus.)
| Métrique | Type | Émise depuis |
|---|---|---|
reindex_time_s | Histogram | _do_reindex fin |
last_reindex_time | Gauge | _do_reindex fin |
db_size | Gauge | _do_reindex fin |
Aucune métrique sur les échecs. Un reindex qui plante incrémente _yoyo_failed_reindex_total... non, ça n'existe pas. Pour détecter les échecs en prod, surveiller last_reindex_time qui ne bouge pas.
Tests
Cf. tests/test_reindex.py. Tests présents :
| Test | Couvre |
|---|---|
test_reindex_inserts_travels | Le pipeline complet ingère les records de la source mockée |
test_reindex_updates_existing_travels | Le ON CONFLICT DO UPDATE met à jour raw_data |
test_reindex_removes_stale_travels | Le DELETE final purge les ids absents de l'amont |
test_reindex_fetches_both_sources | Les deux URLs (base et /seaside) sont hit |
test_seaside_records_tagged | Les records de /seaside ont bien is_seaside=TRUE |
test_reindex_aborts_on_seaside_failure_keeps_old_data | Si /seaside plante, le DELETE final n'est pas exécuté |
test_trailing_slash_in_slow_search_url | SLOW_SEARCH_URL avec ou sans / final produit la bonne URL /seaside |
test_schedule_runner_exists | Smoke test du module schedule_runner |
Pour les tests synchrones, on appelle _do_reindex(pool) directement (pas via le Thread de reindex()).
Points d'extension
Ajouter une source amont
(ex: un autre sous-endpoint /promo)
- Dans
_do_reindex, ajouter un troisième_ingest_source(...). - Décider de l'ordre : la dernière source gagne en cas de conflit d'id.
- Si la source impose un tag (comme
isSeaside), passermark_seaside=Trueou créer un mécanisme génériquemark_with={...}.
Ajouter un champ tagué à l'ingestion
Modèle de la 3.2.0 : isSeaside n'existe pas côté amont, c'est better-search qui l'ajoute. Pour répliquer :
- Modifier
_ingest_sourcepour injecter le tag. - Ajouter une migration qui crée la colonne générée correspondante (
is_xxx BOOLEAN GENERATED ALWAYS AS …). - Ajouter un index si filtre attendu.
- Ajouter le param query côté
server.pyet le WHERE côtésearch_query.py. - Tests.
Optimiser le coût d'embedding
Aujourd'hui, chaque reindex re-embed tous les records (~quelques centaines × 2880 appels/jour). Pour éviter de re-embedder un record dont le texte source n'a pas changé :
- Ajouter une colonne
content_hash TEXT GENERATED ALWAYS AS (md5(get_embedding_text_sql(raw_data))) STORED(ou équivalent). - Modifier
_ingest_sourcepour récupérer les hashes existants avant d'appelerinfomaniak.get_embeddings; ne batcher que les nouveaux/modifiés. - Compute du
get_embedding_text_sqlcôté Postgres ≠ Python — attention à la cohérence (BeautifulSoup vsregexp_replace).
Non implémenté car le coût total est négligeable pour le moment.
Passer en multi-réplique
Si on déploie 2+ pods :
- Remplacer le
threading.Lockpar un advisory lock Postgres :pythonconn.execute("SELECT pg_try_advisory_lock(%s)", (REINDEX_LOCK_ID,)) - Ou déplacer le reindex vers un container
workerdédié, et garder les pods API « read-only ».
Voir ADR 0005.
Logging
Verbeux à INFO. Exemples :
INFO reindex Starting reindex
INFO reindex Loaded 10 records on page 1 of https://horizon.buchard.ch/api/travels
INFO reindex Loaded 10 records on page 2 of https://horizon.buchard.ch/api/travels
...
INFO reindex Loaded 5 records on page 3 of https://horizon.buchard.ch/api/travels/seaside
INFO reindex Removing stale entries...
INFO reindex Reindex finishedLogger "reindex". Pour avoir des logs structurés, ajouter un handler JSON dans app.py.

