From 7669a79f4fa93bfc6d63b2d5c04dc7e618f34cae Mon Sep 17 00:00:00 2001 From: Kazeia Team Date: Thu, 11 Jun 2026 16:10:29 +0200 Subject: [PATCH] =?UTF-8?q?fix(rag):=20ingestion=20atomique=20+=20build/te?= =?UTF-8?q?ardown=20s=C3=A9rialis=C3=A9s=20(code-review)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- .../app/src/main/java/com/kazeia/rag/Rag.kt | 17 ++++++++---- .../app/src/main/java/com/kazeia/rag/RagDb.kt | 27 +++++++++++++++++++ .../java/com/kazeia/service/KazeiaService.kt | 16 ++++++++--- 3 files changed, 51 insertions(+), 9 deletions(-) diff --git a/kazeia-android/app/src/main/java/com/kazeia/rag/Rag.kt b/kazeia-android/app/src/main/java/com/kazeia/rag/Rag.kt index a3eec3e..35e6a51 100644 --- a/kazeia-android/app/src/main/java/com/kazeia/rag/Rag.kt +++ b/kazeia-android/app/src/main/java/com/kazeia/rag/Rag.kt @@ -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(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 } /** (Ré)indexe les documents qui n'ont pas encore d'embeddings. Appelé sur diff --git a/kazeia-android/app/src/main/java/com/kazeia/rag/RagDb.kt b/kazeia-android/app/src/main/java/com/kazeia/rag/RagDb.kt index 27afcbc..630be2b 100644 --- a/kazeia-android/app/src/main/java/com/kazeia/rag/RagDb.kt +++ b/kazeia-android/app/src/main/java/com/kazeia/rag/RagDb.kt @@ -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, 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 { val out = ArrayList() diff --git a/kazeia-android/app/src/main/java/com/kazeia/service/KazeiaService.kt b/kazeia-android/app/src/main/java/com/kazeia/service/KazeiaService.kt index 7fd5006..209c4f9 100644 --- a/kazeia-android/app/src/main/java/com/kazeia/service/KazeiaService.kt +++ b/kazeia-android/app/src/main/java/com/kazeia/service/KazeiaService.kt @@ -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