#!/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). """ from __future__ import annotations import json import math import sys import time import requests import config 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. " "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 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) -> str: top_k = config.get_int("TOP_K", 3) context_chunks = retrieve(index, 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 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')}", 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", "") if text.startswith(REPLY_PREFIX): continue # ecoul propriului raspuns in self-chat, ignorat sender = msg.get("from") print(f"[consumer] {sender}: {text[:80]}", file=sys.stderr) react_seen(sender, msg.get("id"), msg.get("fromMe", False)) send_reply(sender, "Caut informatia, revin imediat...") try: reply = ask_llm(index, text) 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()