#!/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` pentru de ce, si de ce cautarea in index foloseste doar liniile de eroare, nu toata captura. """ from __future__ import annotations import json import math import os import sys import time import requests import config import ocr 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. " "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)." ) def bridge_url() -> str: return f"http://{config.get('BRIDGE_HOST')}:{config.get('BRIDGE_PORT')}" def load_index() -> list[dict]: try: return json.loads(config.INDEX_FILE.read_text(encoding="utf-8")) except OSError: return [] 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 retrieve(index: list[dict], question: str, top_k: int) -> list[str]: if not index: return [] q_vec = embed(question) scored = [(cosine(q_vec, e["embedding"]), e["text"]) for e in index] scored.sort(key=lambda pair: pair[0], reverse=True) return [text for _, text in scored[:top_k]] def ask_llm(index: list[dict], question: str, search_query: str | None = None) -> str: """`search_query` separa CE se cauta in index de CE se trimite modelului. La o captura de ecran, cautarea merge pe liniile de eroare, iar modelului ii dam fereastra intreaga: contextul din jurul erorii ajuta raspunsul, dar strica regasirea. """ top_k = config.get_int("TOP_K", 3) context_chunks = retrieve(index, search_query or question, top_k) context = "\n\n---\n\n".join(context_chunks) if context_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 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 index, 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. """ 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." ) path = media.get("path") try: ocr_text = ocr.run(path) except Exception as exc: # noqa: BLE001 print(f"[consumer] OCR esuat pe {path}: {exc}", file=sys.stderr) ocr_text = "" finally: # Capturile sunt de unica folosinta: pot contine date de client si nu # avem niciun motiv sa le pastram pe disc dupa ce am citit textul. try: if path: os.unlink(path) except OSError: pass if not ocr_text: if text: return text, text, None return None, None, NO_TEXT_IN_IMAGE print(f"[consumer] OCR: {len(ocr_text)} caractere din captura", file=sys.stderr) return build_image_question(text, ocr_text), ocr.retrieval_query(text, ocr_text), None 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)'}", 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() messages = resp.json().get("messages", []) for msg in messages: if msg.get("isGroup"): continue text = msg.get("text", "") or "" if text.startswith(REPLY_PREFIX): continue # ecoul propriului raspuns in self-chat, ignorat sender = msg.get("from") has_image = bool(msg.get("media")) print( f"[consumer] {sender}: {'[imagine] ' if has_image else ''}{text[:80]}", file=sys.stderr, ) react_seen(sender, msg.get("id"), msg.get("fromMe", False)) 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) continue if not question: continue try: reply = ask_llm(index, question, search_query) 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) except Exception as exc: # noqa: BLE001 print(f"[consumer] poll error: {exc}", file=sys.stderr) time.sleep(poll_s) if __name__ == "__main__": main()