Files
ROMFASTSQL/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/consumer.py
Claude Agent af04d97de9 fix(maria): „am trimis la suport" doar cand puntea chiar a trimis, si un log care nu minte
Gasite citind logul rularii de test de azi.

1. `requests` NU ridica exceptie la 4xx/5xx, iar puntea raspunde 503 cand nu e
   conectata la WhatsApp si 500 cand `sendMessage` cade. send_reply/send_image
   ignorau codul, deci escaladarea se inregistra `notified: true` si omul primea
   „te contacteaza cineva" pentru un mesaj care nu plecase nicaieri — exact
   promisiunea pentru care exista ESCALATED_RECORDED. Acum trimiterea intoarce
   motivul esecului ("" la reusita), iar `notified` si `notify_error` vin de acolo.

2. Puntea nu loga nimic la trimitere: o escaladare nu lasa nicio urma pe partea de
   WhatsApp, deci nu se poate verifica daca captura chiar a ajuns la suport.
   /send si /send-image logheaza acum destinatarul si inceputul mesajului.

3. Mesajele primite se logau taiate la 80 de caractere, fara semn ca sunt taiate.
   Un mesaj de exact 80 arata ca unul intreg — asa am ajuns azi la concluzia
   gresita ca gardul de „mesaj prea vag" nu functioneaza, cand de fapt mesajul era
   mai lung decat parea. 200 de caractere si „… (+N)".

4. La o captura pe un fir deschis se logau doua linii „caut dupa" diferite, iar
   prima nu era interogarea folosita. Prima zice acum „din captura, retin".

Teste: 103 pass (una noua: puntea respinge cu 503 -> notified false).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Q4uzvgm7AyJch5WH8QHRhY
2026-09-01 19:44:24 +00:00

672 lines
28 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).
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 fir as fir_mod
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."
)
TACERE_OK = (
'Am inteles, ma opresc. Scrie-mi "Maria, continua" cand vrei sa reiau.'
)
REVENIRE_OK = "Sunt aici. Cu ce te ajut?"
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, istoric: list[dict] | None = None) -> str:
"""Raspunsul modelului. `istoric` = schimburile de pana acum, la o continuare.
Fara istoric, „si acum ce fac?" ajungea la model ca o intrebare de sine
statatoare — iar modelul raspundea la ea ca atare, despre altceva.
"""
context = "\n\n---\n\n".join(chunks) if chunks else "(fara documente indexate)"
user_message = f"CONTEXT:\n{context}\n\nINTREBARE:\n{question}"
mesaje = [{"role": "system", "content": SYSTEM_PROMPT}]
# Istoricul precede contextul: ultimul mesaj trebuie sa fie intrebarea curenta,
# altfel modelul raspunde la penultima.
mesaje += (istoric or [])[:-1]
mesaje.append({"role": "user", "content": user_message})
resp = requests.post(
f"{config.get('LLM_URL')}/v1/chat/completions",
json={
"messages": mesaje,
"max_tokens": config.get_int("MAX_TOKENS", 250),
# Fara asta llama.cpp raspunde cu temperatura lui implicita (0,8).
# Aici nu vrem creativitate: raspunsul trebuie sa fie ce scrie in
# context, la fel de fiecare data. Vezi raspunsurile inventate din
# 2026-09-01 (grup „Maria Test").
"temperature": float(config.get("LLM_TEMPERATURE", "0")),
},
timeout=60,
)
resp.raise_for_status()
return resp.json()["choices"][0]["message"]["content"]
def _trimite(url: str, corp: dict, catre: str, timeout: int) -> str:
"""Trimite prin punte. Intoarce "" daca a plecat, altfel motivul esecului.
`requests` NU ridica exceptie la 4xx/5xx, iar puntea raspunde 503 cand nu e
conectata la WhatsApp si 500 cand `sendMessage` cade. Fara verificarea
codului, escaladarea se inregistra ca „notificata" si Maria promitea „te
contacteaza cineva" pentru un mesaj care nu plecase nicaieri — exact
promisiunea pe care ESCALATED_RECORDED exista ca sa n-o facem.
"""
try:
resp = requests.post(url, json=corp, timeout=timeout)
except Exception as exc: # noqa: BLE001
print(f"[consumer] trimitere esuata catre {catre}: {exc}", file=sys.stderr)
return str(exc)
if resp.status_code >= 400:
motiv = f"HTTP {resp.status_code}: {resp.text[:150]}"
print(f"[consumer] trimitere esuata catre {catre}: {motiv}", file=sys.stderr)
return motiv
return ""
def send_reply(to: str, text: str) -> str:
return _trimite(f"{bridge_url()}/send", {"to": to, "text": REPLY_PREFIX + text}, to, 15)
def send_image(to: str, path: str, caption: str) -> str:
return _trimite(f"{bridge_url()}/send-image",
{"to": to, "path": path, "caption": caption}, to, 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)
# „din captura", nu „caut dupa": pe un fir, interogarea finala e alta (ancora
# plus asta) si se logheaza separat — doua linii „caut dupa" induceau in eroare.
print(f"[consumer] din captura, retin: {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, fir: dict | None = None) -> 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,
}
if fir and fir.get("schimburi"):
record["fir"] = fir["schimburi"]
jid = config.get("SUPPORT_JID") or ""
if jid:
cine = msg.get("pushName") or msg.get("from") or "necunoscut"
discutie = fir_mod.rezumat(fir) if fir else ""
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]}"
+ (f"\n\n--- discutia de pana acum ---\n{discutie}" if discutie else "")
)
if media.get("path") and os.path.exists(media["path"]):
eroare = send_image(jid, media["path"], rezumat[:1000])
else:
eroare = _trimite(f"{bridge_url()}/send", {"to": jid, "text": rezumat}, jid, 20)
record["notified"] = not eroare
if eroare:
record["notify_error"] = eroare
print(f"[consumer] escaladare: notificarea a esuat: {eroare}", 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
REAMINTIRE_MIN = 15 # cat de rar poate Maria sa reaminteasca aceeasi problema
def _fisier_escaladare(ref: str):
d = config.STATE_DIR / "escalations"
if not d.exists():
return None
return next(iter(sorted(d.glob(f"*-{ref}.json"))), None)
def _notifica_suport(text: str, imagine: str | None = None) -> bool:
"""Mesaj catre suport; cu `imagine`, chiar captura primita, nu doar textul citit."""
jid = config.get("SUPPORT_JID") or ""
if not jid:
return False
if imagine and os.path.exists(imagine):
return not send_image(jid, imagine, text[:1000])
return not _trimite(f"{bridge_url()}/send", {"to": jid, "text": text}, jid, 20)
def completeaza(ref: str, msg: dict, text: str) -> str | None:
"""Adauga mesajul la escaladarea deschisa si intoarce ce i se raspunde omului.
`None` daca escaladarea nu mai exista (atunci mesajul se trateaza normal).
Trei feluri de completare, cu trei raspunsuri diferite: un semnal de urgenta
schimba escaladarea si anunta echipa; o intrebare de stare primeste ce stim
(de cat timp asteapta, daca am reamintit); un detaliu primeste o confirmare
scurta, rotita, ca sa nu sune a robot la a treia oara.
"""
fisier = _fisier_escaladare(ref)
if fisier is None:
return None
try:
record = json.loads(fisier.read_text(encoding="utf-8"))
except (OSError, ValueError):
return None
fel = triaj.fel_completare(text)
completari = record.setdefault("completari", [])
completari.append({"at": time.time(), "text": text, "fel": fel,
"message_id": msg.get("id")})
cine = msg.get("pushName") or msg.get("from") or "necunoscut"
minute = int((time.time() - record.get("at", time.time())) / 60)
if fel == "urgenta":
record["urgenta"] = True
trimis = _notifica_suport(
f"[Maria] URGENT — {ref} (de la {cine})\n{text[:800]}")
record["ultima_reamintire"] = time.time()
raspuns = triaj.URGENTA_MARCATA.format(ref=ref)
elif fel == "stare":
# Reamintirea e limitata: la fiecare „tot nimic" nu suna telefonul echipei.
de_reamintit = time.time() - record.get("ultima_reamintire", record.get("at", 0)) \
> REAMINTIRE_MIN * 60
reamintit = False
if de_reamintit:
reamintit = _notifica_suport(
f"[Maria] {ref}: clientul intreaba de {minute} minute daca s-a "
f"rezolvat.\n{text[:400]}")
if reamintit:
record["ultima_reamintire"] = time.time()
intrebari = sum(1 for c in completari if c.get("fel") == "stare")
raspuns = triaj.raspuns_stare(ref, minute, reamintit,
config.get("SUPPORT_PHONE", ""),
a_cata=intrebari - 1)
else:
_notifica_suport(f"[Maria] Completare la {ref} (de la {cine}):\n{text[:800]}",
imagine=(msg.get("media") or {}).get("path"))
raspuns = triaj.confirmare(ref, len(completari) - 1)
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)
print(f"[consumer] completare '{fel}' la {ref}: {text[:60]!r}", file=sys.stderr)
return raspuns
_GRUPURI_CACHE: dict = {"at": 0.0, "date": {}}
GRUPURI_TTL_S = 600
def grup_info(jid: str) -> dict:
"""Cate persoane sunt in grup. Raspunsul se tine 10 minute in memorie."""
if time.time() - _GRUPURI_CACHE["at"] > GRUPURI_TTL_S:
try:
resp = requests.get(f"{bridge_url()}/groups", timeout=15)
resp.raise_for_status()
_GRUPURI_CACHE["date"] = {g["jid"]: g for g in resp.json().get("groups", [])}
_GRUPURI_CACHE["at"] = time.time()
except Exception as exc: # noqa: BLE001
print(f"[consumer] nu pot citi grupurile: {exc}", file=sys.stderr)
return _GRUPURI_CACHE["date"].get(jid, {})
def preluare_de_om(msg: dict) -> bool:
"""Un om din echipa a scris in discutie, deci Maria se retrage.
Doar in grupuri cu mai multi oameni. In self-chat si in grupul de test (un
singur participant) TOT ce se scrie e `fromMe`: acolo regula ar face Maria sa
amuteasca la primul mesaj, deci ramane doar comanda explicita.
"""
if not msg.get("fromMe") or not msg.get("isGroup"):
return False
if (msg.get("text") or "").startswith(REPLY_PREFIX):
return False # propriul raspuns, nu un om
minim = config.get_int("FIR_PRELUARE_MIN_PARTICIPANTI", 2)
return grup_info(msg.get("from") or "").get("participants", 0) >= minim
def grupuri_permise() -> set[str]:
"""JID-urile de grup in care Maria are voie sa raspunda (ALLOWED_GROUP_JIDS)."""
brut = config.get("ALLOWED_GROUP_JIDS", "") or ""
return {j.strip() for j in brut.replace(";", ",").split(",") if j.strip()}
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)
fir = fir_mod.incarca(sender)
# Comenzi explicite. Merg oriunde, si acolo unde regula automata de mai jos nu
# se aplica (self-chat, grup de test cu un singur om).
cmd = fir_mod.comanda(text)
if cmd and not has_image:
fir = fir or fir_mod.deschide(sender, "")
if cmd == "stop":
fir_mod.taci(fir)
fir_mod.salveaza(fir)
print(f"[consumer] {sender}: tac la cerere", file=sys.stderr)
send_reply(sender, TACERE_OK)
else:
fir_mod.vorbeste(fir)
fir_mod.salveaza(fir)
send_reply(sender, REVENIRE_OK)
return
# Un om din echipa a intrat in discutie: Maria nu vorbeste peste el.
if preluare_de_om(msg):
fir = fir or fir_mod.deschide(sender, "")
fir_mod.adauga(fir, "om", text)
fir_mod.taci(fir)
fir_mod.salveaza(fir)
print(f"[consumer] {sender}: a preluat un om, tac {fir_mod.tacere_s() // 60} min",
file=sys.stderr)
return
if fir_mod.tace(fir):
print(f"[consumer] {sender}: fir preluat de om, nu raspund", file=sys.stderr)
return
react_seen(sender, msg.get("id"), msg.get("fromMe", False))
continuare = fir_mod.este_continuare(fir, text, has_image)
# Firul are o escaladare deschisa: mesajul e completare la ea, nu intrebare
# noua. Inaintea confirmarii si a cautarii — raspunsul vine instant, deci un
# „caut informatia, revin imediat" ar fi o promisiune inutila, iar embedding-ul
# s-ar calcula degeaba.
if continuare and fir.get("ref") and text.strip() and not has_image:
raspuns = completeaza(fir["ref"], msg, text.strip())
if raspuns:
fir_mod.adauga(fir, "client", text)
fir_mod.adauga(fir, "maria", raspuns)
fir_mod.salveaza(fir)
send_reply(sender, raspuns)
return
fir["ref"] = None # escaladarea nu mai exista; tratam mesajul normal
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
if continuare:
fir_mod.adauga(fir, "client", text or "[captura de ecran]")
fir_mod.extinde_ancora(fir, msg.get("ocr_text") or "")
# Captura pe o escaladare deschisa e chiar detaliul care lipsea, nu o
# problema noua: pleaca la aceeasi referinta, cu tot cu imagine.
if fir.get("ref") and has_image:
raspuns = completeaza(fir["ref"], msg, msg.get("ocr_text") or text)
if raspuns:
fir_mod.adauga(fir, "maria", raspuns)
fir_mod.salveaza(fir)
send_reply(sender, raspuns)
return
fir["ref"] = None # escaladarea nu mai exista; tratam mesajul normal
search_query = fir_mod.interogare(fir, search_query)
print(f"[consumer] continuare pe firul deschis; caut dupa {search_query[:90]!r}",
file=sys.stderr)
else:
fir = fir_mod.deschide(sender, msg.get("ocr_text") or text)
fir_mod.adauga(fir, "client", text or "[captura de ecran]")
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"]:
# „Am o eroare la o factura", fara sa spuna care: nu avem ce cauta si nici
# ce trimite la suport — programatorul ar pune exact aceeasi intrebare.
# O singura data pe fir: daca nici dupa ce a fost intrebat nu spune mai
# mult, mesajul pleaca la om, care poate insista altfel.
if triaj.prea_vag(text, has_image) and not fir.get("detalii_cerute"):
fir["detalii_cerute"] = True
fir_mod.adauga(fir, "maria", triaj.CERE_DETALII)
fir_mod.salveaza(fir)
print("[consumer] mesaj prea vag -> cer detalii", file=sys.stderr)
send_reply(sender, triaj.CERE_DETALII)
return
print(f"[consumer] fara acoperire in documente -> suport ({result['reason']})", file=sys.stderr)
record = escalate(msg, question, result, fir)
sablon = ESCALATED_SENT if record["notified"] else ESCALATED_RECORDED
raspuns = sablon.format(ref=record["ref"])
fir["ref"] = record["ref"]
fir_mod.adauga(fir, "maria", raspuns)
fir_mod.salveaza(fir)
send_reply(sender, raspuns)
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, fir)
raspuns = triaj.raspuns_escaladat(
result["chunks"][0], nivel, record["ref"], record["notified"], cu_imagine=has_image)
fir["ref"] = record["ref"]
fir_mod.adauga(fir, "maria", raspuns)
fir_mod.salveaza(fir)
send_reply(sender, raspuns)
return
try:
reply = ask_llm(result["chunks"], question, fir_mod.istoric(fir) if continuare else None)
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."
fir_mod.adauga(fir, "maria", reply)
fir_mod.salveaza(fir)
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") and msg.get("from") not in grupuri_permise():
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()