Skip to content

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

python
# 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."""
python
# src.schedule_runner
def start_schedule():
    """Lance un thread qui boucle sur reindex() toutes REINDEX_FREQUENCY secondes."""
python
# src.db
class BuchardDatabase:
    def reindex(self):
        """Threaded reindex job — proxy vers reindex.reindex(self._pool)."""
python
# 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

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


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_running est lu/écrit sous le lock pour éviter une race entre with _reindex_lock et Thread(...).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 à True ad vitam.

Pipeline d'ingestion (_do_reindex)

Cf. architecture/reindex-lifecycle.md pour le détail. Résumé :

python
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..n jusqu'à totalFilteredRecords == 0.
  • Pour chaque page :
    • Si mark_seaside, injecte record['isSeaside'] = True (pour la colonne générée is_seaside).
    • Calcule les embeddings via infomaniak.get_embeddings(texts) (batch).
    • INSERT … ON CONFLICT (id) DO UPDATE sur chaque record.
    • Commit par page.
  • 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 : DELETE n'est atteint que si les deux sources se sont bien terminées.

Configuration

Env varDéfautEffet
REINDEX_FREQUENCY900Secondes entre cycles auto (entier strict)
SLOW_SEARCH_URLobligatoireBase URL de l'amont, sans trailing / requis

Constantes en code :

ConstanteValeurEffet
_FETCH_TIMEOUT30 sTimeout HTTP requests.get
Page size amont10Hardcodé dans le ?size=10 de l'URL

Métriques émises

(Voir modules/monitoring.md pour le détail Prometheus.)

MétriqueTypeÉmise depuis
reindex_time_sHistogram_do_reindex fin
last_reindex_timeGauge_do_reindex fin
db_sizeGauge_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 :

TestCouvre
test_reindex_inserts_travelsLe pipeline complet ingère les records de la source mockée
test_reindex_updates_existing_travelsLe ON CONFLICT DO UPDATE met à jour raw_data
test_reindex_removes_stale_travelsLe DELETE final purge les ids absents de l'amont
test_reindex_fetches_both_sourcesLes deux URLs (base et /seaside) sont hit
test_seaside_records_taggedLes records de /seaside ont bien is_seaside=TRUE
test_reindex_aborts_on_seaside_failure_keeps_old_dataSi /seaside plante, le DELETE final n'est pas exécuté
test_trailing_slash_in_slow_search_urlSLOW_SEARCH_URL avec ou sans / final produit la bonne URL /seaside
test_schedule_runner_existsSmoke 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)

  1. Dans _do_reindex, ajouter un troisième _ingest_source(...).
  2. Décider de l'ordre : la dernière source gagne en cas de conflit d'id.
  3. Si la source impose un tag (comme isSeaside), passer mark_seaside=True ou créer un mécanisme générique mark_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 :

  1. Modifier _ingest_source pour injecter le tag.
  2. Ajouter une migration qui crée la colonne générée correspondante (is_xxx BOOLEAN GENERATED ALWAYS AS …).
  3. Ajouter un index si filtre attendu.
  4. Ajouter le param query côté server.py et le WHERE côté search_query.py.
  5. 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é :

  1. Ajouter une colonne content_hash TEXT GENERATED ALWAYS AS (md5(get_embedding_text_sql(raw_data))) STORED (ou équivalent).
  2. Modifier _ingest_source pour récupérer les hashes existants avant d'appeler infomaniak.get_embeddings ; ne batcher que les nouveaux/modifiés.
  3. Compute du get_embedding_text_sql côté Postgres ≠ Python — attention à la cohérence (BeautifulSoup vs regexp_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 :

  1. Remplacer le threading.Lock par un advisory lock Postgres :
    python
    conn.execute("SELECT pg_try_advisory_lock(%s)", (REINDEX_LOCK_ID,))
  2. Ou déplacer le reindex vers un container worker dé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 finished

Logger "reindex". Pour avoir des logs structurés, ajouter un handler JSON dans app.py.

Contributors

No contributors

Changelog

No recent changes