diff --git a/scripts/rag_ingest_resume.py b/scripts/rag_ingest_resume.py new file mode 100644 index 0000000..04f4bdd --- /dev/null +++ b/scripts/rag_ingest_resume.py @@ -0,0 +1,174 @@ +#!/usr/bin/env python3 +""" +Ingest JSONL into ChromaDB with resume support. +- Прерывается по Ctrl+C безопасно +- При перезапуске продолжает с места останова +- Не удаляет ничего без спроса +""" +import json, os, sys, time, urllib.request, glob, signal +from chromadb import PersistentClient + +EF_URL = "http://localhost:8081/embedding" +DB_DIR = os.path.expanduser("~/nubes/chroma_db") +STATE_FILE = os.path.join(DB_DIR, ".ingest_state.txt") + +ALL_DIRS = [ + "~/nubes/data/elsa_jsonl_wi_en", + "~/nubes/data/elsa_mdb_jsonl", + "~/nubes/data/elsa_mdb_jsonl/ipsvrap", + "~/nubes/data/elsa_jsonl_htm_en", + "~/nubes/data/elsa_jsonl", + "~/nubes/data/elsa_mdb_jsonl/dbsvrfi", + "~/nubes/data/elsa_mdb_jsonl/dbsvrfz", +] + +running = True + +def handle_sigint(sig, frame): + global running + print("\n⏳ Graceful stop after current file...") + running = False + +def embed(texts): + for attempt in range(5): + try: + data = json.dumps({"content": texts}).encode() + req = urllib.request.Request(EF_URL, data=data, + headers={"Content-Type": "application/json"}) + with urllib.request.urlopen(req, timeout=120) as resp: + body = json.loads(resp.read()) + if isinstance(body, list): + return [item["embedding"][0] for item in body] + return body.get("embedding", []) + except Exception as e: + if attempt < 4: + sys.stdout.write(f"R{attempt}") + sys.stdout.flush() + time.sleep(15) + else: + raise + +def load_state(): + if os.path.exists(STATE_FILE): + with open(STATE_FILE) as f: + return set(line.strip() for line in f if line.strip()) + return set() + +def save_state(fpath): + os.makedirs(os.path.dirname(STATE_FILE), exist_ok=True) + with open(STATE_FILE, "a") as f: + f.write(fpath + "\n") + +def main(): + signal.signal(signal.SIGINT, handle_sigint) + + processed = load_state() + + # Подключаемся к БД + client = PersistentClient(path=DB_DIR) + existing = os.path.exists(os.path.join(DB_DIR, "chroma.sqlite3")) + + if existing: + try: + collection = client.get_collection("elsa_docs") + print(f"DB exists: {collection.count()} docs") + except: + collection = client.create_collection("elsa_docs") + print("Created new collection") + else: + collection = client.create_collection("elsa_docs") + print("Created new ChromaDB") + + # Собираем все файлы + all_files = [] + for d in ALL_DIRS: + path = os.path.expanduser(d) + files = sorted(glob.glob(os.path.join(path, "*.jsonl"))) + all_files.extend(files) + + pending = [f for f in all_files if f not in processed] + print(f"Files: total={len(all_files)}, done={len(processed)}, pending={len(pending)}") + + if not pending: + print("All done.") + return + + total = collection.count() + errors = 0 + + for fpath in pending: + if not running: + print("Stopped by user") + break + + docs, ids = [], [] + prefix = os.path.basename(os.path.dirname(fpath)) + "_" + fname = os.path.basename(fpath) + + try: + with open(fpath) as fh: + for line in fh: + try: + d = json.loads(line) + title = d.get("title", "") or d.get("full_title", "") + text = d.get("text", "") + doc = (title + "\n" + text).strip()[:1000] + if len(doc) < 20: + continue + docs.append(doc) + ids.append(f"{prefix}{fname}_{len(docs)}") + except: + continue + except Exception as e: + print(f" [ERR] {fname}: {e}") + errors += 1 + save_state(fpath) + continue + + if not docs: + save_state(fpath) + print(f" {fname}: empty") + continue + + sys.stdout.write(f" {fname}: {len(docs)} docs") + sys.stdout.flush() + + ok = True + for j in range(0, len(docs), 50): + if not running: + ok = False + break + batch_docs = docs[j:j+50] + batch_ids = ids[j:j+50] + for retry in range(5): + try: + embs = embed(batch_docs) + collection.add(ids=batch_ids, documents=batch_docs, + embeddings=embs, metadatas=[{"source": prefix.strip("_")}] * len(batch_docs)) + total += len(batch_docs) + sys.stdout.write(".") + sys.stdout.flush() + break + except Exception as e: + if retry < 4: + sys.stdout.write(f"R{retry}") + sys.stdout.flush() + time.sleep(20) + else: + sys.stdout.write("X") + sys.stdout.flush() + errors += len(batch_docs) + + time.sleep(0.2) + + if ok: + save_state(fpath) + + print(f" total={collection.count()}") + + print(f"\nTotal: {collection.count()} docs. Errors: {errors}") + if pending and not running: + print("Restart to continue from where stopped.") + +if __name__ == "__main__": + main()