#!/usr/bin/env python3 """Consuma mesajele din puntea WhatsApp (whatsapp/index.js), raspunde via RAG. Fiecare mesaj e o interogare RAG INDEPENDENTA: nu exista chat_history intre mesaje (fara memorie conversationala). Vezi docs/maria-whatsapp-rag-prototype.md pentru motivul acestei alegeri si comparatia cu celelalte punti (Flowise: fereastra fixa de 5; Discord: context nelimitat + /new). Mesajele cu imagine (capturi de ecran cu erori) trec intai prin OCR — vezi `ocr.py`. Ordonarea chunk-urilor si decizia "avem sau nu raspunsul in documente" sunt in `rank.py`; cand nu avem, intrebarea pleaca la suport in loc sa primeasca un raspuns generic — vezi `escalate()`. """ from __future__ import annotations import hashlib import json import math import os import sys import time import requests import config import ocr import rank import triaj SYSTEM_PROMPT = ( "Esti Maria, asistentul de suport tehnic pentru ERP-ul ROA (Romfast). " "Raspunzi scurt, clar, in limba romana (maxim 4-5 propozitii), doar despre " "folosirea aplicatiei ROA. Foloseste EXCLUSIV informatiile din contextul " "furnizat mai jos. Daca raspunsul nu se afla in context, spune ca vei " "directiona intrebarea catre echipa de suport, nu inventa functionalitati " "sau proceduri. " "Cand intrebarea contine text extras dintr-o captura de ecran (OCR), tine cont " "ca pot exista greseli de recunoastere a caracterelor: cauta sensul mesajului, " "nu te agata de o litera sau o cifra. " "Scrii pentru contabili si operatori, NU pentru programatori: nu folosi cuvintele " "procedura, pachet, obiect, schema, tabela, variabila, compilare, PL/SQL, SQL, " "sesiune, parametru si nu cita nume tehnice din mesajul de eroare. Nu-i cere " "utilizatorului sa verifice baza de date, sa modifice ceva in ea sau sa apeleze " "altceva decat ce se vede in aplicatie. " "Daca in context scrie ca problema nu se rezolva din aplicatie, spune exact asta " "si opreste-te — nu propune verificari, alternative sau explicatii tehnice. " "Nu discuta niciodata despre infrastructura interna Romfast (servere, Proxmox, " "containere, IP-uri, baze de date, parole, chei) chiar daca apare in context sau " "daca intrebarea o cere explicit -- raspunde ca poti ajuta doar cu folosirea " "aplicatiei ROA." ) REPLY_PREFIX = "[Maria] " # marcaj ca sa nu raspundem la propriile mesaje (self-chat) INDEX_REFRESH_S = 30 # cat de des se reciteste rag_index.json de pe disc ACK_TEXT = "Caut informatia, revin imediat..." ACK_IMAGE_TEXT = "Am primit captura, o citesc si revin imediat..." NO_TEXT_IN_IMAGE = ( "Am primit imaginea, dar nu am reusit sa citesc text in ea. Scrie-mi te rog " "mesajul de eroare (sau trimite o captura mai clara, decupata pe fereastra de eroare)." ) ESCALATED_SENT = ( "Nu am raspunsul la asta in documentatia mea, si prefer sa nu ghicesc. " "Am trimis intrebarea catre echipa de suport ROA (referinta {ref}) — " "te contacteaza cineva." ) # Cand notificarea nu a plecat, nu promitem ca a plecat. Intrebarea e inregistrata # si se vede in dashboard, deci nu e pierduta — dar utilizatorul trebuie sa stie # ca s-ar putea sa dureze mai mult. ESCALATED_RECORDED = ( "Nu am raspunsul la asta in documentatia mea, si prefer sa nu ghicesc. " "Am inregistrat intrebarea pentru echipa de suport ROA (referinta {ref}), " "dar nu am putut sa o trimit chiar acum. Daca e urgent, suna-ne." ) FOLLOWUP_OK = ( "Multumesc, am adaugat si asta la {ref} — echipa vede detaliul cand preia " "problema." ) def bridge_url() -> str: return f"http://{config.get('BRIDGE_HOST')}:{config.get('BRIDGE_PORT')}" class Index: """Chunk-urile indexate, plus statisticile lexicale pentru BM25. BM25 se reconstruieste la fiecare recitire a indexului. La ordinul de marime de aici (sute de chunk-uri) costa milisecunde, si e mai bine decat un al doilea fisier pe disc care se poate desincroniza de rag_index.json. """ def __init__(self, entries: list[dict]): self.entries = entries self.texts = [e["text"] for e in entries] self.sources = [e.get("source", "?") for e in entries] self.bm25 = rank.Bm25(self.texts) def __len__(self) -> int: return len(self.entries) def load_index() -> Index: try: return Index(json.loads(config.INDEX_FILE.read_text(encoding="utf-8"))) except OSError: return Index([]) def embed(text: str) -> list[float]: resp = requests.post( f"{config.get('OLLAMA_URL')}/api/embeddings", json={"model": config.get("EMBED_MODEL"), "prompt": text}, timeout=30, ) resp.raise_for_status() return resp.json()["embedding"] def cosine(a: list[float], b: list[float]) -> float: dot = sum(x * y for x, y in zip(a, b)) na = math.sqrt(sum(x * x for x in a)) nb = math.sqrt(sum(y * y for y in b)) return dot / (na * nb) if na and nb else 0.0 def search(index: Index, query: str, top_k: int) -> dict: """Cauta si decide daca avem acoperire. Intoarce si dovezile, pentru log.""" if not len(index): return {"chunks": [], "covered": False, "reason": "index gol", "best_cosine": 0.0, "top": []} q_vec = embed(query) cosines = [cosine(q_vec, e["embedding"]) for e in index.entries] bm_scores = index.bm25.scores(query) order = rank.rank(cosines, bm_scores, rank.by_codes(query, index.texts))[:top_k] texts = [index.texts[i] for i in order] # cosinusul chunk-urilor CHIAR date modelului, nu cel mai bun din tot indexul: # altfel acoperirea se judeca pe o dovada care nu ajunge in context. best = max((cosines[i] for i in order), default=0.0) covered, reason = rank.assess(query, best, texts, index.bm25) return { "chunks": texts, "covered": covered, "reason": reason, "best_cosine": best, "top": [ {"source": index.sources[i], "cosine": round(cosines[i], 3), "bm25": round(bm_scores[i], 2)} for i in order ], } def ask_llm(chunks: list[str], question: str) -> str: context = "\n\n---\n\n".join(chunks) if chunks else "(fara documente indexate)" user_message = f"CONTEXT:\n{context}\n\nINTREBARE:\n{question}" resp = requests.post( f"{config.get('LLM_URL')}/v1/chat/completions", json={ "messages": [ {"role": "system", "content": SYSTEM_PROMPT}, {"role": "user", "content": user_message}, ], "max_tokens": config.get_int("MAX_TOKENS", 250), }, timeout=60, ) resp.raise_for_status() return resp.json()["choices"][0]["message"]["content"] def send_reply(to: str, text: str) -> None: requests.post(f"{bridge_url()}/send", json={"to": to, "text": REPLY_PREFIX + text}, timeout=15) def send_image(to: str, path: str, caption: str) -> None: requests.post( f"{bridge_url()}/send-image", json={"to": to, "path": path, "caption": caption}, timeout=60, ) def react_seen(to: str, message_id: str, from_me: bool) -> None: try: requests.post( f"{bridge_url()}/react", json={"to": to, "id": message_id, "emoji": "\U0001F440", "fromMe": from_me}, timeout=10, ) except Exception as exc: # noqa: BLE001 print(f"[consumer] react error: {exc}", file=sys.stderr) def build_image_question(caption: str, ocr_text: str) -> str: """Ce vede modelul cand mesajul a fost o captura de ecran.""" parts = [] if caption.strip(): parts.append(caption.strip()) parts.append( "Utilizatorul a trimis o captura de ecran. Text extras automat din imagine " f"(OCR, poate contine greseli de recunoastere):\n---\n{ocr_text}\n---" ) if not caption.strip(): parts.append("Explica-i ce inseamna eroarea si cum o rezolva.") return "\n\n".join(parts) def prepare_query(msg: dict) -> tuple[str | None, str | None, str | None]: """(intrebare pentru model, interogare pentru cautare, raspuns imediat de trimis). Al treilea element e diferit de None cand nu se poate raspunde deloc (imaginea nu s-a descarcat, OCR indisponibil, nimic lizibil in captura) — atunci textul lui pleaca asa cum e, fara sa mai deranjam modelul. Captura NU se sterge aici: mai e nevoie de ea daca intrebarea pleaca la suport. Stergerea o face apelantul, dupa ce mesajul e complet tratat. """ text = (msg.get("text") or "").strip() media = msg.get("media") if not media: return text, text, None if media.get("error"): print(f"[consumer] imagine nedescarcata: {media['error']}", file=sys.stderr) if text: return text, text, None return None, None, ( "Nu am reusit sa descarc imaginea. Mai incearca o data, sau scrie-mi " "mesajul de eroare ca text." ) try: ocr_text = ocr.run(media.get("path")) except Exception as exc: # noqa: BLE001 print(f"[consumer] OCR esuat pe {media.get('path')}: {exc}", file=sys.stderr) ocr_text = "" msg["ocr_text"] = ocr_text if not ocr_text: print("[consumer] OCR: nimic lizibil in captura", file=sys.stderr) if text: return text, text, None return None, None, NO_TEXT_IN_IMAGE # Textul extras se logheaza INTEGRAL. Fara asta, un raspuns gresit nu se poate # explica: nu stii daca a citit prost captura sau a cautat prost in documente. print( f"[consumer] OCR ({len(ocr_text)} caractere):\n" + "\n".join(" | " + line for line in ocr_text.splitlines()), file=sys.stderr, ) cautare = ocr.retrieval_query(text, ocr_text) print(f"[consumer] caut dupa: {cautare!r}", file=sys.stderr) return build_image_question(text, ocr_text), cautare, None def reference(when: float, message_id: str | None) -> str: """Referinta scurta, pe care omul o poate cita la telefon: M-260831-A1B2. Se calculeaza din id-ul mesajului, ca sa fie stabila daca se reia escaladarea aceluiasi mesaj, si sa nu depinda de un contor pastrat undeva. """ seed = (message_id or str(when)).encode("utf-8", "replace") sufix = hashlib.sha1(seed).hexdigest()[:4].upper() return f"M-{time.strftime('%y%m%d', time.localtime(when))}-{sufix}" def escalate(msg: dict, question: str, search_result: dict) -> dict: """Trimite intrebarea la suport si o inregistreaza in jurnal. Intoarce inregistrarea. Jurnalul se scrie INTOTDEAUNA, si cand notificarea esueaza sau `SUPPORT_JID` nu e configurat — altfel o intrebare fara raspuns dispare fara urma, exact cazul pe care escaladarea trebuia sa-l rezolve. Cine cheama functia afla din `notified` daca poate promite utilizatorului ca mesajul chiar a plecat. """ media = msg.get("media") or {} now = time.time() record = { "ref": reference(now, msg.get("id")), "at": now, "from": msg.get("from"), "push_name": msg.get("pushName"), "message_id": msg.get("id"), "text": msg.get("text"), "ocr_text": msg.get("ocr_text"), "question": question, "reason": search_result.get("reason"), "best_cosine": round(search_result.get("best_cosine", 0.0), 3), "top": search_result.get("top"), "had_image": bool(media.get("path")), "notified": False, } jid = config.get("SUPPORT_JID") or "" if jid: cine = msg.get("pushName") or msg.get("from") or "necunoscut" rezumat = ( f"[Maria] Intrebare fara raspuns in documente — {record['ref']}\n" f"De la: {cine}\n" f"Motiv: {record['reason']}\n\n" f"{(msg.get('ocr_text') or msg.get('text') or '').strip()[:1200]}" ) try: if media.get("path") and os.path.exists(media["path"]): send_image(jid, media["path"], rezumat[:1000]) else: requests.post(f"{bridge_url()}/send", json={"to": jid, "text": rezumat}, timeout=20) record["notified"] = True except Exception as exc: # noqa: BLE001 record["notify_error"] = str(exc) print(f"[consumer] escaladare: notificarea a esuat: {exc}", file=sys.stderr) else: print("[consumer] escaladare: SUPPORT_JID nesetat — doar in jurnal", file=sys.stderr) try: d = config.STATE_DIR / "escalations" d.mkdir(parents=True, exist_ok=True) nume = f"{int(record['at'])}-{record['ref']}.json" (d / nume).write_text(json.dumps(record, ensure_ascii=False, indent=2), encoding="utf-8") print(f"[consumer] escaladare {record['ref']}: {nume} (notificat: {record['notified']})", file=sys.stderr) except OSError as exc: print(f"[consumer] escaladare: nu pot scrie jurnalul: {exc}", file=sys.stderr) return record PENDING_FILE_NAME = "pending.json" PENDING_TTL_S = 30 * 60 # cat timp un raspuns scurt mai e „completare", nu intrebare noua def _pending_path(): return config.STATE_DIR / "escalations" / PENDING_FILE_NAME def _pending_all() -> dict: try: return json.loads(_pending_path().read_text(encoding="utf-8")) except (OSError, ValueError): return {} def _pending_write(data: dict) -> None: try: _pending_path().parent.mkdir(parents=True, exist_ok=True) _pending_path().write_text(json.dumps(data, ensure_ascii=False), encoding="utf-8") except OSError as exc: print(f"[consumer] nu pot scrie {PENDING_FILE_NAME}: {exc}", file=sys.stderr) def pending_set(sender: str, ref: str) -> None: """Retine ca l-am intrebat pe om cat e de urgent, ca sa stiu unde duce raspunsul.""" data = _pending_all() data[sender] = {"ref": ref, "at": time.time()} _pending_write(data) def pending_get(sender: str) -> dict | None: intrare = _pending_all().get(sender) if not intrare or time.time() - intrare.get("at", 0) > PENDING_TTL_S: return None return intrare def pending_clear(sender: str) -> None: data = _pending_all() if data.pop(sender, None) is not None: _pending_write(data) def append_followup(ref: str, msg: dict, text: str) -> bool: """Adauga raspunsul omului la escaladarea deschisa si il trimite la suport. Raspunsul la „te blocheaza sau poti continua?" e chiar informatia care lipseste din escaladare. O escaladare noua ar rupe firul: aceeasi problema, alta referinta, iar omul a citat-o deja pe prima. """ d = config.STATE_DIR / "escalations" fisier = next((f for f in sorted(d.glob(f"*-{ref}.json"))), None) if d.exists() else None if fisier is None: return False try: record = json.loads(fisier.read_text(encoding="utf-8")) except (OSError, ValueError): return False record.setdefault("completari", []).append( {"at": time.time(), "text": text, "message_id": msg.get("id")}) try: fisier.write_text(json.dumps(record, ensure_ascii=False, indent=2), encoding="utf-8") except OSError as exc: print(f"[consumer] nu pot actualiza {fisier.name}: {exc}", file=sys.stderr) jid = config.get("SUPPORT_JID") or "" if jid: cine = msg.get("pushName") or msg.get("from") or "necunoscut" try: requests.post(f"{bridge_url()}/send", json={ "to": jid, "text": f"[Maria] Completare la {ref} (de la {cine}):\n{text[:800]}", }, timeout=20) except Exception as exc: # noqa: BLE001 print(f"[consumer] completare netrimisa: {exc}", file=sys.stderr) print(f"[consumer] completare la {ref}: {text[:60]!r}", file=sys.stderr) return True def cleanup_media(msg: dict) -> None: """Capturile pot contine date de client — nu raman pe disc dupa ce s-a tratat mesajul.""" path = (msg.get("media") or {}).get("path") if not path: return try: os.unlink(path) except OSError: pass def handle(index: Index, msg: dict) -> None: sender = msg.get("from") has_image = bool(msg.get("media")) text = msg.get("text", "") or "" print(f"[consumer] {sender}: {'[imagine] ' if has_image else ''}{text[:80]}", file=sys.stderr) react_seen(sender, msg.get("id"), msg.get("fromMe", False)) # Raspunsul la „te blocheaza sau poti continua?" merge la escaladarea deschisa. # Inainte de orice altceva: nu e o intrebare noua, deci nu se cauta si nu se # confirma cu „caut informatia". asteptat = pending_get(sender) if not has_image else None if asteptat and text.strip() and not rank.codes(text): if append_followup(asteptat["ref"], msg, text.strip()): pending_clear(sender) send_reply(sender, FOLLOWUP_OK.format(ref=asteptat["ref"])) return pending_clear(sender) # „Am o eroare", fara sa spuna care: nu ghicim si nu deranjam suportul — intrebam. if triaj.prea_vag(text, has_image): print("[consumer] mesaj prea vag -> cer detalii", file=sys.stderr) send_reply(sender, triaj.CERE_DETALII) return send_reply(sender, ACK_IMAGE_TEXT if has_image else ACK_TEXT) question, search_query, immediate = prepare_query(msg) if immediate is not None: send_reply(sender, immediate) return if not question: return try: result = search(index, search_query, config.get_int("TOP_K", 3)) except Exception as exc: # noqa: BLE001 print(f"[consumer] cautare esuata: {exc}", file=sys.stderr) send_reply(sender, "Scuze, am o problema tehnica momentan. Cineva din echipa te va contacta.") return surse = ", ".join(f"{t['source']}({t['cosine']})" for t in result["top"]) or "-" print(f"[consumer] rank: {surse} | {result['reason']}", file=sys.stderr) if not result["covered"]: print(f"[consumer] fara acoperire in documente -> suport ({result['reason']})", file=sys.stderr) record = escalate(msg, question, result) sablon = ESCALATED_SENT if record["notified"] else ESCALATED_RECORDED send_reply(sender, sablon.format(ref=record["ref"])) return # Eroare care nu se rezolva din aplicatie: raspunsul se compune din dictionar # si pleaca automat la programatori. Vezi triaj.py pentru de ce nu prin model. nivel = triaj.nivel_suport(result["chunks"][0]) if result["chunks"] else None if nivel: print(f"[consumer] eroare de nivel '{nivel}' -> raspuns fix + suport", file=sys.stderr) record = escalate(msg, question, result) send_reply(sender, triaj.raspuns_escaladat( result["chunks"][0], nivel, record["ref"], record["notified"], cu_imagine=has_image)) pending_set(sender, record["ref"]) return try: reply = ask_llm(result["chunks"], question) except Exception as exc: # noqa: BLE001 print(f"[consumer] LLM error: {exc}", file=sys.stderr) reply = "Scuze, am o problema tehnica momentan. Cineva din echipa te va contacta." send_reply(sender, reply) print(f"[consumer] -> raspuns catre {sender}", file=sys.stderr) def main() -> None: poll_s = config.get_int("POLL_INTERVAL_S", 2) print( f"[consumer] polling {bridge_url()}/messages la {poll_s}s, " f"LLM={config.get('LLM_URL')}, RAG top-{config.get('TOP_K')}, " f"OCR={'da' if ocr.available() else 'INDISPONIBIL (tesseract lipseste)'}, " f"suport={config.get('SUPPORT_JID') or 'NESETAT (doar jurnal)'}", file=sys.stderr, ) index = load_index() last_index_check = time.time() while True: try: if time.time() - last_index_check > INDEX_REFRESH_S: index = load_index() last_index_check = time.time() resp = requests.get(f"{bridge_url()}/messages", timeout=10) resp.raise_for_status() for msg in resp.json().get("messages", []): if msg.get("isGroup"): continue if (msg.get("text") or "").startswith(REPLY_PREFIX): continue # ecoul propriului raspuns in self-chat, ignorat try: handle(index, msg) finally: cleanup_media(msg) except Exception as exc: # noqa: BLE001 print(f"[consumer] poll error: {exc}", file=sys.stderr) time.sleep(poll_s) if __name__ == "__main__": main()