Doua greseli din grupul "Maria Test", 2026-09-01, pe care 27409fb le-a atenuat dar
nu le-a rezolvat.
1. La o reclamatie vaga Maria cauta, in loc sa intrebe. "Buna . Am si eu o factura
pe luna august cu eroare in spv la trimitere AUTO SULE" nu spune CE eroare e —
n-are ce cauta in documente si n-are ce trimite la suport, fiindca programatorul
ar pune exact aceeasi intrebare. Gardul (triaj.prea_vag) exista, dar cerea text
sub 12 cuvinte; mesajul are 14. Gresea si invers: "nu pot incarca factura in SPV,
imi da eroare de certificat" are 9 cuvinte si ARE raspuns in documente, si era
oprita degeaba. Lungimea nu masoara cat de precis e mesajul.
-> gardul nu mai numara cuvinte si se aplica DUPA cautare, doar cand nu exista
acoperire. O singura data pe fir: daca nici intrebat omul nu spune mai mult,
mesajul pleaca la suport.
2. Textul si captura, trimise una dupa alta, erau tratate ca doua probleme fara
legatura. `este_continuare` rupea firul la ORICE imagine — dar cazul frecvent e
tocmai omul care scrie problema si trimite captura imediat dupa, sau care
raspunde la "trimite-mi o captura". Textul se pierdea, captura se cauta doar pe
OCR-ul ei, se deschideau doua fire si se puteau deschide doua escaladari pentru
aceeasi problema.
-> o captura in primele FIR_IMAGINE_MIN (5) minute continua firul: textul citit
din ea intra in ancora (cu tot cu coduri), cautarea se face pe mesaj +
captura, iar pe o escaladare deschisa pleaca la aceeasi referinta — cu
imaginea, nu doar cu textul citit din ea (_notifica_suport ia si media).
Si, ca urmare a lui (2): RANK_STRONG_COSINE dispare. Interogarea pe un fir e ancora
plus mesajul nou, deci cosinusul urca la fiecare replica fara sa apara vreo dovada
noua — aceeasi captura da 0,778 singura si 0,836 cu mesajul de dinainte, adica peste
0,80 pus ieri. Orice prag fix de sus e trecut de o discutie destul de lunga. Acoperirea
ramane pe dovada lexicala (termenii distinctivi ai intrebarii chiar in chunk-uri),
care e stabila la lungime. Niciun caz cu raspuns in documente nu avea nevoie de
scurtatura.
ops/calibrate-rank.py: 26/26, cu interogarea combinata adaugata la set. Teste: 102 pass.
Verificat end-to-end pe scenariul real (mesaj, apoi captura la 10 secunde): intrebare
de detalii, apoi o singura escaladare care poarta si textul si captura.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Q4uzvgm7AyJch5WH8QHRhY
659 lines
27 KiB
Python
659 lines
27 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 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, 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 "")
|
|
)
|
|
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
|
|
|
|
|
|
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
|
|
try:
|
|
if imagine and os.path.exists(imagine):
|
|
send_image(jid, imagine, text[:1000])
|
|
else:
|
|
requests.post(f"{bridge_url()}/send", json={"to": jid, "text": text}, timeout=20)
|
|
return True
|
|
except Exception as exc: # noqa: BLE001
|
|
print(f"[consumer] mesaj catre suport netrimis: {exc}", file=sys.stderr)
|
|
return False
|
|
|
|
|
|
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()
|