Files
ROMFASTSQL/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/consumer.py
Claude Agent c5af5f2380 feat(maria): interzice explicit discutiile despre infrastructura interna
Maria raspunde clientilor ERP ROA. Depozitul ei de documente contine doar
material orientat spre client final, dar adaugam si o bariera in SYSTEM_PROMPT:
refuza intrebarile despre servere, Proxmox, containere, IP-uri, baze de date,
parole sau chei, chiar daca ceva de acest fel ajunge in context.

Regula completa (ce are voie sa intre in ~/.maria-bridge/documents/) e scrisa
in /workspace/claude-agent/CLAUDE.md, sectiunea "Maria (RAG) - GRANITA
OBLIGATORIE".

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Q4uzvgm7AyJch5WH8QHRhY
2026-08-31 17:38:51 +00:00

152 lines
5.4 KiB
Python

#!/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()