diff --git a/ingest/mao_ingest.py b/ingest/mao_ingest.py new file mode 100644 index 0000000..0ec41c0 --- /dev/null +++ b/ingest/mao_ingest.py @@ -0,0 +1,363 @@ +""" +Ingestion MAO PDFs → ChromaDB +Collection : mao_docs +Embedding : nomic-embed-text (Ollama local) +Chunking structurel : paragraphes + sections, ~450 mots, overlap 80 mots + +Double usage : + 1. RAG assistant — répond aux questions sur les VSTs/plugins + 2. Accélération wiki — fournit les chunks pertinents au LLM pour générer les pages wiki + (à la place de lire le PDF entier) +""" +import json, os, re, subprocess, tempfile, time, urllib.request, urllib.parse +from pathlib import Path + +PDFFETCH_URL = "http://localhost:15001/files" +OLLAMA_EMBED = "http://localhost:11434/api/embeddings" +CHROMA_URL = "http://localhost:8800" +CHROMA_BASE = "/api/v2/tenants/default_tenant/databases/default_database" +COLLECTION = "mao_docs" +EMBED_MODEL = "nomic-embed-text" +STATE_FILE = Path("/home/ubuntu/mao/ingest_state.json") + +CHUNK_WORDS = 450 +OVERLAP_WORDS = 80 +MIN_WORDS = 40 + +# Correspondances nom de fichier → métadonnées produit +# Enrichit la recherche : "comment fonctionne la reverb de Serum ?" → filtre type=synth, product=Serum +PRODUCT_HINTS = { + "serum": {"product": "Serum", "manufacturer": "Xfer", "type": "synth"}, + "massive": {"product": "Massive", "manufacturer": "Native Instruments", "type": "synth"}, + "vital": {"product": "Vital", "manufacturer": "Matt Tytel", "type": "synth"}, + "surge": {"product": "Surge XT", "manufacturer": "Surge Synth", "type": "synth"}, + "kontakt": {"product": "Kontakt", "manufacturer": "Native Instruments", "type": "sampler"}, + "battery": {"product": "Battery", "manufacturer": "Native Instruments", "type": "sampler"}, + "ableton": {"product": "Ableton Live", "manufacturer": "Ableton", "type": "daw"}, + "live": {"product": "Ableton Live", "manufacturer": "Ableton", "type": "daw"}, + "fl_studio": {"product": "FL Studio", "manufacturer": "Image-Line", "type": "daw"}, + "fl studio": {"product": "FL Studio", "manufacturer": "Image-Line", "type": "daw"}, + "logic": {"product": "Logic Pro", "manufacturer": "Apple", "type": "daw"}, + "cubase": {"product": "Cubase", "manufacturer": "Steinberg", "type": "daw"}, + "pro_tools": {"product": "Pro Tools", "manufacturer": "Avid", "type": "daw"}, + "reaper": {"product": "Reaper", "manufacturer": "Cockos", "type": "daw"}, + "izotope": {"product": "iZotope", "manufacturer": "iZotope", "type": "mastering"}, + "ozone": {"product": "Ozone", "manufacturer": "iZotope", "type": "mastering"}, + "neutron": {"product": "Neutron", "manufacturer": "iZotope", "type": "mixing"}, + "fabfilter": {"product": "FabFilter", "manufacturer": "FabFilter", "type": "effects"}, + "pro-q": {"product": "Pro-Q", "manufacturer": "FabFilter", "type": "eq"}, + "pro-r": {"product": "Pro-R", "manufacturer": "FabFilter", "type": "reverb"}, + "pro-c": {"product": "Pro-C", "manufacturer": "FabFilter", "type": "compressor"}, + "valhalla": {"product": "Valhalla", "manufacturer": "Valhalla DSP", "type": "reverb"}, + "waves": {"product": "Waves", "manufacturer": "Waves", "type": "effects"}, + "u-he": {"product": "u-he", "manufacturer": "u-he", "type": "synth"}, + "diva": {"product": "Diva", "manufacturer": "u-he", "type": "synth"}, + "hive": {"product": "Hive", "manufacturer": "u-he", "type": "synth"}, + "repro": {"product": "Repro", "manufacturer": "u-he", "type": "synth"}, + "zebra": {"product": "Zebra", "manufacturer": "u-he", "type": "synth"}, + "spire": {"product": "Spire", "manufacturer": "Reveal Sound", "type": "synth"}, + "sylenth": {"product": "Sylenth1", "manufacturer": "LennarDigital","type": "synth"}, + "omnisphere": {"product": "Omnisphere", "manufacturer": "Spectrasonics","type": "synth"}, + "nexus": {"product": "Nexus", "manufacturer": "refx", "type": "rompler"}, + "reaktor": {"product": "Reaktor", "manufacturer": "Native Instruments", "type": "modular"}, + "max_msp": {"product": "Max/MSP", "manufacturer": "Cycling '74", "type": "modular"}, + "soundtoys": {"product": "SoundToys", "manufacturer": "SoundToys", "type": "effects"}, + "rc-20": {"product": "RC-20", "manufacturer": "XLN Audio", "type": "effects"}, + "decimort": {"product": "Decimort", "manufacturer": "D16 Group", "type": "effects"}, + "d16": {"product": "D16 Group", "manufacturer": "D16 Group", "type": "effects"}, +} + + +# ── État ────────────────────────────────────────────────────────────────────── + +def load_state(): + if STATE_FILE.exists(): + return json.loads(STATE_FILE.read_text()) + return {"ingested": [], "failed": []} + +def save_state(s): + STATE_FILE.parent.mkdir(parents=True, exist_ok=True) + STATE_FILE.write_text(json.dumps(s, indent=2)) + + +# ── ChromaDB ────────────────────────────────────────────────────────────────── + +def chroma_request(method, path, data=None): + url = CHROMA_URL + path + body = json.dumps(data).encode() if data else None + req = urllib.request.Request(url, data=body, method=method, + headers={"Content-Type": "application/json"}) + with urllib.request.urlopen(req, timeout=30) as r: + return json.loads(r.read()) + +def get_or_create_collection(): + try: + r = chroma_request("POST", f"{CHROMA_BASE}/collections", { + "name": COLLECTION, + "metadata": {"hnsw:space": "cosine"} + }) + print(f"Collection créée : {COLLECTION} ({r['id']})") + return r["id"] + except urllib.error.HTTPError as e: + if e.code == 409: + cols = chroma_request("GET", f"{CHROMA_BASE}/collections") + col_id = next(c["id"] for c in cols if c["name"] == COLLECTION) + print(f"Collection existante : {COLLECTION} ({col_id})") + return col_id + raise + + +def add_documents(col_id, ids, embeddings, documents, metadatas): + chroma_request("POST", f"{CHROMA_BASE}/collections/{col_id}/add", { + "ids": ids, "embeddings": embeddings, + "documents": documents, "metadatas": metadatas + }) + + +# ── Embeddings ──────────────────────────────────────────────────────────────── + +def embed(text): + payload = json.dumps({"model": EMBED_MODEL, "prompt": text}).encode() + req = urllib.request.Request(OLLAMA_EMBED, data=payload, + headers={"Content-Type": "application/json"}, method="POST") + with urllib.request.urlopen(req, timeout=60) as r: + return json.loads(r.read())["embedding"] + + +# ── PDF helpers ─────────────────────────────────────────────────────────────── + +def list_mao_pdfs(): + req = urllib.request.Request(PDFFETCH_URL, + headers={"User-Agent": "mao-ingest/1.0"}) + with urllib.request.urlopen(req, timeout=15) as r: + all_files = json.loads(r.read()) + # PDFs MAO = racine (pas de sous-dossier "fallout/") + return [f for f in all_files if not f.startswith("fallout/") and f.endswith(".pdf")] + +def download_pdf(filename): + url = f"{PDFFETCH_URL}/{urllib.parse.quote(filename)}" + req = urllib.request.Request(url, headers={"User-Agent": "mao-ingest/1.0"}) + with urllib.request.urlopen(req, timeout=300) as r: + return r.read() + +def extract_text(pdf_bytes): + with tempfile.NamedTemporaryFile(suffix=".pdf", delete=False) as f: + f.write(pdf_bytes); tmp = f.name + text = "" + try: + r = subprocess.run(["pdftotext", "-layout", tmp, "-"], + capture_output=True, timeout=180) + if r.returncode == 0: + text = r.stdout.decode("utf-8", errors="replace") + finally: + os.unlink(tmp) + if len(text.strip()) < 200: + try: + import io + from pypdf import PdfReader + reader = PdfReader(io.BytesIO(pdf_bytes)) + text = "\n".join(p.extract_text() or "" for p in reader.pages) + except Exception as e: + raise RuntimeError(f"Extraction impossible: {e}") + return text + + +# ── Chunking structurel ─────────────────────────────────────────────────────── + +def is_section_header(line): + l = line.strip() + if not l or len(l) > 100: + return False + if l.isupper() and 2 <= len(l.split()) <= 10: + return True + if re.match(r'^(Chapter|Section|Part|Chapitre|Partie|\d+[\.\)])\s+\S', l, re.I): + return True + return False + +def chunk_by_structure(text): + text = re.sub(r'[ \t]+', ' ', text) + text = re.sub(r'\n{3,}', '\n\n', text) + + raw_blocks = re.split(r'\n{2,}', text) + + blocks = [] + i = 0 + while i < len(raw_blocks): + b = raw_blocks[i].strip() + if not b: + i += 1; continue + if is_section_header(b) and i + 1 < len(raw_blocks): + merged = b + "\n\n" + raw_blocks[i+1].strip() + blocks.append(merged) + i += 2 + else: + blocks.append(b) + i += 1 + + chunks = [] + current_blocks = [] + current_words = 0 + + for block in blocks: + bwords = len(block.split()) + if bwords < 5: + continue + if current_words + bwords > CHUNK_WORDS and current_blocks: + chunk = "\n\n".join(current_blocks).strip() + if len(chunk.split()) >= MIN_WORDS: + chunks.append(chunk) + overlap_blocks = [] + overlap_words = 0 + for ob in reversed(current_blocks): + ow = len(ob.split()) + if overlap_words + ow <= OVERLAP_WORDS: + overlap_blocks.insert(0, ob) + overlap_words += ow + else: + break + current_blocks = overlap_blocks + [block] + current_words = overlap_words + bwords + else: + current_blocks.append(block) + current_words += bwords + + if current_blocks: + chunk = "\n\n".join(current_blocks).strip() + if len(chunk.split()) >= MIN_WORDS: + chunks.append(chunk) + + return chunks + + +# ── Métadonnées ─────────────────────────────────────────────────────────────── + +def infer_metadata(filename): + f = filename.lower() + name = Path(filename).name + stem = Path(filename).stem.lower().replace("-", "_").replace(" ", "_") + + meta = { + "source": name, + "path": filename, + "product": "unknown", + "manufacturer": "unknown", + "type": "unknown", + "language": "fr" if any(k in f for k in ["_fr", "francais", "manuel"]) else "en", + } + + for hint_key, hint_vals in PRODUCT_HINTS.items(): + if hint_key in f: + meta.update(hint_vals) + break + + # Fallback : nom de fichier sans extension comme product + if meta["product"] == "unknown": + meta["product"] = Path(filename).stem[:40] + + return meta + + +# ── Query helper (pour wiki pipeline) ──────────────────────────────────────── + +def query_for_wiki(col_id, product_name, top_k=15): + """ + Retrouve les chunks les plus pertinents pour générer une page wiki. + Utiliser avec : query = "features workflow controls {product_name}" + """ + query = f"features workflow interface controls parameters {product_name}" + emb = embed(query) + result = chroma_request("POST", f"{CHROMA_BASE}/collections/{col_id}/query", { + "query_embeddings": [emb], + "n_results": top_k, + "where": {"product": product_name}, + "include": ["documents", "metadatas", "distances"] + }) + return result + + +# ── Main ────────────────────────────────────────────────────────────────────── + +def main(): + import sys + # Mode : ingest (défaut) ou query + mode = sys.argv[1] if len(sys.argv) > 1 else "ingest" + + print(f"=== MAO ChromaDB — mode: {mode} ===") + print(f"Collection : {COLLECTION}") + print(f"Modèle embedding : {EMBED_MODEL}\n") + + state = load_state() + col_id = get_or_create_collection() + + if mode == "query": + product = sys.argv[2] if len(sys.argv) > 2 else "Serum" + print(f"Query wiki pour : {product}") + r = query_for_wiki(col_id, product) + for i, (doc, meta) in enumerate(zip(r["documents"][0], r["metadatas"][0])): + print(f"\n--- Chunk {i+1} [{meta['source']}] ---") + print(doc[:300] + "...") + return + + pdfs = list_mao_pdfs() + print(f"{len(pdfs)} PDFs MAO sur Guardian\n") + + for filename in pdfs: + if filename in state["ingested"]: + print(f"[SKIP] {filename}") + continue + + print(f"\n[{filename}]") + try: + pdf_bytes = download_pdf(filename) + print(f" {len(pdf_bytes)/1024:.0f} Ko téléchargés") + + text = extract_text(pdf_bytes) + print(f" {len(text)} chars extraits") + + chunks = chunk_by_structure(text) + meta_base = infer_metadata(filename) + print(f" {len(chunks)} chunks | {meta_base['product']} ({meta_base['type']}) | lang:{meta_base['language']}") + + ids, embeddings, documents, metadatas = [], [], [], [] + + for i, chunk in enumerate(chunks): + chunk_id = f"{Path(filename).stem}_{i:04d}" + t0 = time.time() + emb = embed(chunk) + word_count = len(chunk.split()) + ids.append(chunk_id) + embeddings.append(emb) + documents.append(chunk) + metadatas.append({ + **meta_base, + "chunk_index": i, + "total_chunks": len(chunks), + "word_count": word_count, + }) + print(f" chunk {i+1}/{len(chunks)} ({word_count}w) OK ({time.time()-t0:.1f}s)") + + if len(ids) >= 50: + add_documents(col_id, ids, embeddings, documents, metadatas) + ids, embeddings, documents, metadatas = [], [], [], [] + + if ids: + add_documents(col_id, ids, embeddings, documents, metadatas) + + print(f"\n Indexé — {len(chunks)} chunks") + state["ingested"].append(filename) + state["failed"] = [f for f in state["failed"] if f.get("file") != filename] + save_state(state) + + except Exception as e: + print(f" ERREUR: {e}") + state["failed"].append({"file": filename, "error": str(e)}) + save_state(state) + + done = len(state["ingested"]) + failed = len(state["failed"]) + print(f"\n=== Terminé : {done} PDFs indexés, {failed} échecs ===") + for f in state.get("failed", []): + print(f" ECHEC: {f['file']} — {f['error']}") + +if __name__ == "__main__": + main()