fix(rag): ingestion atomique + build/teardown sérialisés (code-review)

Rag.ingest : embed D'ABORD, n'écrit la base qu'ensuite, via la nouvelle
RagDb.replaceSource (delete+insert dans UNE transaction). Un échec transitoire
d'embedding ne détruit plus les chunks existants (auparavant deleteSource était
appelé avant la boucle d'embed → corpus tronqué/vidé en silence, non rattrapé
par reingestMissing qui ne reprend que les docs à zéro chunk).

KazeiaService.buildRag/teardownRag : synchronized(ragLock) + buildRag idempotent
(retourne l'instance existante si déjà active). onCreate et un toggle config RAG
peuvent se chevaucher (deux coroutines serviceScope) → sans verrou, deux
EngineEmbedder natifs étaient créés et le perdant fuyait (jamais close()).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Kazeia Team 2026-06-11 16:10:29 +02:00
parent e1d9c94751
commit 7669a79f4f
3 changed files with 51 additions and 9 deletions

View File

@ -45,19 +45,26 @@ class Rag(
fun ingest(source: String, rawText: String): Int {
db.upsertDoc(source, rawText, System.currentTimeMillis())
if (!embedder.isReady) { log("ingest: embedder non prêt (doc enregistré, embeddings différés)"); return 0 }
db.deleteSource(source)
val chunks = Chunker.chunk(rawText, source)
var n = 0
// Embed D'ABORD, ne touche à la base qu'ensuite : un échec d'embedding ne
// doit jamais détruire les chunks existants (sinon corpus tronqué/vidé en
// silence, car reingestMissing ne reprend QUE les docs à zéro chunk).
val ready = ArrayList<RagDb.Chunk>(chunks.size)
for ((i, ch) in chunks.withIndex()) {
// e5 : les passages doivent porter le préfixe "passage: " (symétrique au
// "query: " de retrieve). EngineEmbedder.embedPassage le gère ; le fallback
// FakeEmbedder embed() sans préfixe (symétrique aussi).
val v = (embedder as? EngineEmbedder)?.embedPassage(ch) ?: embedder.embed(ch) ?: continue
db.insert(source, i, ch, modelTag, v); n++
ready.add(RagDb.Chunk(i, ch, v))
}
log("ingest '$source': $n chunks")
if (chunks.isNotEmpty() && ready.isEmpty()) {
log("ingest '$source': embedding KO pour tous les chunks — anciens chunks conservés")
return 0
}
db.replaceSource(source, ready, modelTag) // delete + insert atomiques
log("ingest '$source': ${ready.size} chunks")
loadIndex()
return n
return ready.size
}
/** ()indexe les documents qui n'ont pas encore d'embeddings. Appelé sur

View File

@ -134,6 +134,33 @@ class RagDb private constructor(context: Context)
return writableDatabase.insert(TABLE, null, cv)
}
data class Chunk(val position: Int, val text: String, val vec: FloatArray)
/**
* Remplace ATOMIQUEMENT les chunks d'une source (suppression + insertions
* dans une seule transaction). Tant que la transaction n'est pas committée,
* les anciens chunks restent visibles un échec en cours de route ne laisse
* jamais le corpus à moitié réindexé. L'appelant ne doit appeler ceci que
* lorsqu'il a déjà calculé les vecteurs (cf Rag.ingest : embed d'abord).
*/
fun replaceSource(source: String, chunks: List<Chunk>, model: String) {
val db = writableDatabase
db.beginTransaction()
try {
db.delete(TABLE, "$COL_SOURCE=?", arrayOf(source))
for (c in chunks) {
val cv = ContentValues().apply {
put(COL_SOURCE, source); put(COL_POS, c.position); put(COL_TEXT, c.text)
put(COL_MODEL, model); put(COL_DIM, c.vec.size); put(COL_VEC, vecToBlob(c.vec))
}
db.insert(TABLE, null, cv)
}
db.setTransactionSuccessful()
} finally {
db.endTransaction()
}
}
/** Charge tous les chunks (id, texte, vecteur) pour peupler l'index RAM. */
fun loadAll(): List<Row> {
val out = ArrayList<Row>()

View File

@ -76,6 +76,10 @@ Pas d'introduction, pas d'explication. Juste les 4 lignes en francais. /no_think
// RAG optionnel (default OFF). Inerte tant que l'embedder kazeia-engine n'est
// pas prêt (embedText absent / modèle d'embedding non présent) -> aucune injection.
@Volatile private var rag: com.kazeia.rag.Rag? = null
// Sérialise build/teardown RAG : onCreate et un toggle config peuvent se
// chevaucher (deux coroutines) → sans verrou, deux EngineEmbedder natifs sont
// créés et le perdant fuit (jamais close()).
private val ragLock = Any()
private lateinit var stt: SttEngine
// Cycle de vie Whisper STT — lazy load + swap intelligent.
@ -2063,8 +2067,11 @@ Pas d'introduction, pas d'explication. Juste les 4 lignes en francais. /no_think
}
/** Construit l'instance RAG (embedder e5 + index) si ragEnabled. Pose RagHolder. */
private fun buildRag(): com.kazeia.rag.Rag? {
if (!runtimeConfig.ragEnabled) { log("[RAG] désactivé (config)"); com.kazeia.rag.RagHolder.instance = null; return null }
private fun buildRag(): com.kazeia.rag.Rag? = synchronized(ragLock) {
if (!runtimeConfig.ragEnabled) { log("[RAG] désactivé (config)"); com.kazeia.rag.RagHolder.instance = null; rag = null; return null }
// Idempotent : si une instance est déjà en place (course onCreate/toggle),
// ne pas en recréer une seconde (sinon 2e contexte natif e5 → fuite).
rag?.let { log("[RAG] déjà actif, build ignoré"); return it }
return try {
val embModel = "${KazeiaApplication.MODELS_DIR}/embed-e5-small.gguf"
val emb = com.kazeia.rag.EngineEmbedder(embModel, nThreads = 4, pooling = -1,
@ -2074,13 +2081,14 @@ Pas d'introduction, pas d'explication. Juste les 4 lignes en francais. /no_think
if (it.isReady && it.count() == 0 && it.listDocs().isEmpty())
it.ingestDir("${KazeiaApplication.MODELS_DIR}/../rag_corpus")
it.reingestMissing() // embed les docs admin ajoutés hors-ligne
rag = it
com.kazeia.rag.RagHolder.instance = it
log("[RAG] activé (ready=${it.isReady}, chunks=${it.count()})")
}
} catch (e: Throwable) { log("[RAG] init échec: ${e.message}"); com.kazeia.rag.RagHolder.instance = null; null }
} catch (e: Throwable) { log("[RAG] init échec: ${e.message}"); com.kazeia.rag.RagHolder.instance = null; rag = null; null }
}
private fun teardownRag() {
private fun teardownRag() = synchronized(ragLock) {
rag?.let { try { it.close() } catch (_: Exception) {} }
rag = null
com.kazeia.rag.RagHolder.instance = null