Add Maria WhatsApp+RAG bridge as a service (LXC 171)
Move the /tmp prototype (Baileys bridge + RAG consumer) into git as a proper sibling project to discord-bridge/: own systemd --user units (whatsapp bridge, rag consumer, dashboard, periodic Drive sync timer), a filesystem document store with a stdlib control dashboard (start/ stop/restart, document CRUD, reindex, Google Drive sync via rclone), and an idempotent ops/install.sh following the same conventions. Co-Authored-By: Claude Agent <noreply@anthropic.com>
This commit is contained in:
106
proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/config.py
Normal file
106
proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/config.py
Normal file
@@ -0,0 +1,106 @@
|
||||
"""Configuratie comuna a puntii WhatsApp+RAG pentru Maria (LXC 171 claude-agent).
|
||||
|
||||
Citeste ~/.maria-bridge/env (KEY=value, tolerant la comentarii si ghilimele).
|
||||
Mirrors deliberat conventiile din discord-bridge/config.py, ca sa fie un singur
|
||||
model de citit pentru cine intretine ambele punti pe acest container.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import pathlib
|
||||
|
||||
# Directorul de baza poate fi mutat in teste prin MARIA_BRIDGE_DIR.
|
||||
_DEFAULT_DIR = pathlib.Path.home() / ".maria-bridge"
|
||||
|
||||
STATE_DIR: pathlib.Path = pathlib.Path(os.environ.get("MARIA_BRIDGE_DIR") or _DEFAULT_DIR)
|
||||
DOCS_DIR: pathlib.Path = STATE_DIR / "documents"
|
||||
INDEX_FILE: pathlib.Path = STATE_DIR / "rag_index.json"
|
||||
LOG_DIR: pathlib.Path = STATE_DIR / "logs"
|
||||
ENV_FILE: pathlib.Path = STATE_DIR / "env"
|
||||
AUTH_DIR: pathlib.Path = STATE_DIR / "whatsapp-auth"
|
||||
|
||||
_env: dict[str, str] = {}
|
||||
|
||||
# Valori implicite. LLM_URL/OLLAMA_URL presupun tunel/proxy local catre backend-ul
|
||||
# real (vezi docs/maria-whatsapp-rag-prototype.md) — de completat in env dupa caz.
|
||||
DEFAULTS: dict[str, str] = {
|
||||
"BRIDGE_HOST": "127.0.0.1",
|
||||
"BRIDGE_PORT": "8099",
|
||||
"LLM_URL": "http://127.0.0.1:8091",
|
||||
"OLLAMA_URL": "http://127.0.0.1:11434",
|
||||
"EMBED_MODEL": "nomic-embed-text",
|
||||
"TOP_K": "3",
|
||||
"MAX_TOKENS": "250",
|
||||
"POLL_INTERVAL_S": "2",
|
||||
"TEST_MODE_SELF_CHAT_ONLY": "true",
|
||||
"DASHBOARD_BIND": "127.0.0.1",
|
||||
"DASHBOARD_PORT": "18792",
|
||||
# Tinta rclone pentru sincronizarea depozitului de documente, ex:
|
||||
# "gdrive:romfast/document_store" (dosarul D:\GoogleDrive\romfast\document_store
|
||||
# de pe Windows, vazut prin Google Drive API). Gol = sincronizare dezactivata,
|
||||
# doar upload manual din dashboard.
|
||||
"DRIVE_REMOTE": "",
|
||||
}
|
||||
|
||||
|
||||
def parse_env(text: str) -> dict[str, str]:
|
||||
"""Parseaza un fisier de tip KEY=value. Nu arunca niciodata."""
|
||||
out: dict[str, str] = {}
|
||||
for raw in text.splitlines():
|
||||
line = raw.strip()
|
||||
if not line or line.startswith("#"):
|
||||
continue
|
||||
if line.startswith("export "):
|
||||
line = line[len("export "):].strip()
|
||||
if "=" not in line:
|
||||
continue
|
||||
key, _, val = line.partition("=")
|
||||
key = key.strip()
|
||||
if not key:
|
||||
continue
|
||||
val = val.strip()
|
||||
if val[:1] not in ("'", '"'):
|
||||
cut = val.find(" #")
|
||||
if cut >= 0:
|
||||
val = val[:cut].rstrip()
|
||||
if len(val) >= 2 and val[0] == val[-1] and val[0] in ("'", '"'):
|
||||
val = val[1:-1]
|
||||
out[key] = val
|
||||
return out
|
||||
|
||||
|
||||
def reload(base_dir: str | os.PathLike | None = None) -> dict[str, str]:
|
||||
"""Recalculeaza caile si reciteste env-ul. Returneaza dictionarul incarcat."""
|
||||
global STATE_DIR, DOCS_DIR, INDEX_FILE, LOG_DIR, ENV_FILE, AUTH_DIR, _env
|
||||
if base_dir is None:
|
||||
base_dir = os.environ.get("MARIA_BRIDGE_DIR") or _DEFAULT_DIR
|
||||
STATE_DIR = pathlib.Path(base_dir)
|
||||
DOCS_DIR = STATE_DIR / "documents"
|
||||
INDEX_FILE = STATE_DIR / "rag_index.json"
|
||||
LOG_DIR = STATE_DIR / "logs"
|
||||
ENV_FILE = STATE_DIR / "env"
|
||||
AUTH_DIR = STATE_DIR / "whatsapp-auth"
|
||||
try:
|
||||
_env = parse_env(ENV_FILE.read_text(encoding="utf-8"))
|
||||
except OSError:
|
||||
_env = {}
|
||||
return _env
|
||||
|
||||
|
||||
def get(key: str, default: str | None = None) -> str | None:
|
||||
if key in _env:
|
||||
return _env[key]
|
||||
if key in DEFAULTS:
|
||||
return DEFAULTS[key]
|
||||
return default
|
||||
|
||||
|
||||
def get_int(key: str, default: int) -> int:
|
||||
try:
|
||||
return int(get(key, str(default)) or default)
|
||||
except (TypeError, ValueError):
|
||||
return default
|
||||
|
||||
|
||||
reload()
|
||||
@@ -0,0 +1,147 @@
|
||||
#!/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."
|
||||
)
|
||||
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()
|
||||
@@ -0,0 +1,63 @@
|
||||
#!/usr/bin/env python3
|
||||
"""(Re)construieste rag_index.json din toate documentele din depozit (store.py),
|
||||
cu embeddings Ollama. Rulat manual sau declansat din dashboard (`/api/reindex`)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
|
||||
import requests
|
||||
|
||||
import config
|
||||
import store
|
||||
|
||||
|
||||
def chunk_text(text: str) -> list[str]:
|
||||
# imparte pe linii goale in paragrafe, uneste bucatile mici cu urmatoarea
|
||||
raw_parts = re.split(r"\n\s*\n", text.strip())
|
||||
chunks: list[str] = []
|
||||
buffer = ""
|
||||
for part in raw_parts:
|
||||
part = part.strip()
|
||||
if not part:
|
||||
continue
|
||||
buffer = f"{buffer}\n\n{part}" if buffer else part
|
||||
if len(buffer) >= 200:
|
||||
chunks.append(buffer)
|
||||
buffer = ""
|
||||
if buffer:
|
||||
chunks.append(buffer)
|
||||
return chunks
|
||||
|
||||
|
||||
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=60,
|
||||
)
|
||||
resp.raise_for_status()
|
||||
return resp.json()["embedding"]
|
||||
|
||||
|
||||
def build() -> dict:
|
||||
entries = []
|
||||
docs = store.list_documents()
|
||||
for doc in docs:
|
||||
text = store.read_document(doc["name"])
|
||||
for i, chunk in enumerate(chunk_text(text)):
|
||||
vec = embed(chunk)
|
||||
entries.append({"source": doc["name"], "chunk": i, "text": chunk, "embedding": vec})
|
||||
config.STATE_DIR.mkdir(parents=True, exist_ok=True)
|
||||
config.INDEX_FILE.write_text(json.dumps(entries, ensure_ascii=False), encoding="utf-8")
|
||||
return {"documents": len(docs), "chunks": len(entries)}
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
result = build()
|
||||
print(
|
||||
f"[indexer] {result['documents']} documente, {result['chunks']} chunk-uri -> {config.INDEX_FILE}",
|
||||
file=sys.stderr,
|
||||
)
|
||||
@@ -0,0 +1,2 @@
|
||||
# Consumer + indexer RAG pentru Maria. Restul e stdlib.
|
||||
requests==2.32.3
|
||||
@@ -0,0 +1,53 @@
|
||||
"""Depozitul de documente pentru indexarea RAG a lui Maria.
|
||||
|
||||
Fisiere text (.txt/.md) sub STATE_DIR/documents/, un singur nivel (fara
|
||||
subdirectoare), ca numele afisat in dashboard sa fie neambiguu si sa poata servi
|
||||
direct ca parametru de request fara riscuri de traversare de cale.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
|
||||
import config
|
||||
|
||||
_SAFE_NAME = re.compile(r"^[A-Za-z0-9._-]{1,200}$")
|
||||
|
||||
|
||||
def validate_name(name: str) -> str:
|
||||
if not name or not _SAFE_NAME.match(name) or ".." in name or "/" in name:
|
||||
raise ValueError(f"nume de document invalid: {name!r}")
|
||||
if not name.endswith((".txt", ".md")):
|
||||
raise ValueError("doar fisiere .txt sau .md")
|
||||
return name
|
||||
|
||||
|
||||
def list_documents() -> list[dict]:
|
||||
config.DOCS_DIR.mkdir(parents=True, exist_ok=True)
|
||||
out = []
|
||||
for f in sorted(config.DOCS_DIR.glob("*")):
|
||||
if not f.is_file() or f.suffix not in (".txt", ".md"):
|
||||
continue
|
||||
st = f.stat()
|
||||
out.append({"name": f.name, "size": st.st_size, "mtime": st.st_mtime})
|
||||
return out
|
||||
|
||||
|
||||
def read_document(name: str) -> str:
|
||||
name = validate_name(name)
|
||||
return (config.DOCS_DIR / name).read_text(encoding="utf-8")
|
||||
|
||||
|
||||
def write_document(name: str, content: str) -> None:
|
||||
name = validate_name(name)
|
||||
config.DOCS_DIR.mkdir(parents=True, exist_ok=True)
|
||||
(config.DOCS_DIR / name).write_text(content, encoding="utf-8")
|
||||
|
||||
|
||||
def delete_document(name: str) -> bool:
|
||||
name = validate_name(name)
|
||||
path = config.DOCS_DIR / name
|
||||
if not path.exists():
|
||||
return False
|
||||
path.unlink()
|
||||
return True
|
||||
@@ -0,0 +1,92 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Sincronizeaza depozitul de documente cu un remote rclone (Google Drive) si
|
||||
reconstruieste indexul RAG DOAR daca s-a schimbat efectiv ceva pe disc.
|
||||
|
||||
De ce prin rclone si nu direct cu Google Drive API: containerul e headless
|
||||
(fara browser pentru OAuth interactiv) — rclone se configureaza o data cu un
|
||||
cont de serviciu (`rclone config`, tip `drive`, `service_account_file=...`),
|
||||
vezi README.md, sectiunea "Sincronizare cu Google Drive".
|
||||
|
||||
Config (`~/.maria-bridge/env`, vezi ops/env.example):
|
||||
DRIVE_REMOTE - tinta rclone, ex: gdrive:romfast/document_store
|
||||
(gol = sincronizare dezactivata, doar upload manual din dashboard)
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
|
||||
import config
|
||||
import indexer
|
||||
|
||||
STATE_FILE_NAME = ".sync_state.json"
|
||||
|
||||
|
||||
def _fingerprint() -> str:
|
||||
"""Amprenta continutului depozitului (nume+mtime+marime), ca sa reindexam
|
||||
doar cand s-a schimbat efectiv ceva, nu la fiecare tur de sincronizare."""
|
||||
config.DOCS_DIR.mkdir(parents=True, exist_ok=True)
|
||||
h = hashlib.sha256()
|
||||
for f in sorted(config.DOCS_DIR.glob("*")):
|
||||
if f.is_file() and f.suffix in (".txt", ".md"):
|
||||
st = f.stat()
|
||||
h.update(f"{f.name}:{st.st_mtime_ns}:{st.st_size}\n".encode())
|
||||
return h.hexdigest()
|
||||
|
||||
|
||||
def _sync_state_path():
|
||||
return config.STATE_DIR / STATE_FILE_NAME
|
||||
|
||||
|
||||
def read_sync_state() -> dict:
|
||||
try:
|
||||
return json.loads(_sync_state_path().read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
return {}
|
||||
|
||||
|
||||
def _write_state(fingerprint: str, extra: dict | None = None) -> None:
|
||||
data = {"fingerprint": fingerprint, "synced_at": time.time()}
|
||||
if extra:
|
||||
data.update(extra)
|
||||
config.STATE_DIR.mkdir(parents=True, exist_ok=True)
|
||||
_sync_state_path().write_text(json.dumps(data), encoding="utf-8")
|
||||
|
||||
|
||||
def pull_from_drive() -> dict:
|
||||
"""`rclone sync <DRIVE_REMOTE> -> DOCS_DIR`. Fara remote configurat, e no-op."""
|
||||
remote = config.get("DRIVE_REMOTE")
|
||||
if not remote:
|
||||
return {"ok": True, "skipped": "DRIVE_REMOTE nesetat in env"}
|
||||
config.DOCS_DIR.mkdir(parents=True, exist_ok=True)
|
||||
try:
|
||||
r = subprocess.run(
|
||||
["rclone", "sync", remote, str(config.DOCS_DIR),
|
||||
"--include", "*.txt", "--include", "*.md"],
|
||||
capture_output=True, text=True, timeout=300,
|
||||
)
|
||||
except FileNotFoundError:
|
||||
return {"ok": False, "error": "rclone nu e instalat — vezi README.md"}
|
||||
except subprocess.TimeoutExpired:
|
||||
return {"ok": False, "error": "rclone a depasit timpul (300s)"}
|
||||
return {"ok": r.returncode == 0, "stdout": r.stdout[-2000:], "stderr": r.stderr[-2000:]}
|
||||
|
||||
|
||||
def sync_and_reindex(force: bool = False) -> dict:
|
||||
pulled = pull_from_drive()
|
||||
fp = _fingerprint()
|
||||
last = read_sync_state().get("fingerprint")
|
||||
if not force and fp == last:
|
||||
return {"pulled": pulled, "reindexed": False, "reason": "fara schimbari"}
|
||||
result = indexer.build()
|
||||
_write_state(fp, {"last_build": result})
|
||||
return {"pulled": pulled, "reindexed": True, "build": result}
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
out = sync_and_reindex(force="--force" in sys.argv[1:])
|
||||
print(json.dumps(out, ensure_ascii=False), file=sys.stderr)
|
||||
Reference in New Issue
Block a user