From dd5553e32754049da841026605a7f1b1164f4ef9 Mon Sep 17 00:00:00 2001 From: Claude Agent Date: Mon, 31 Aug 2026 13:06:15 +0000 Subject: [PATCH] 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 --- .../maria-whatsapp-bridge/.gitignore | 4 + .../maria-whatsapp-bridge/README.md | 140 ++++++ .../maria-whatsapp-bridge/dashboard/api.py | 419 ++++++++++++++++++ .../dashboard/index.html | 211 +++++++++ .../dashboard/login.html | 47 ++ .../maria-whatsapp-bridge/ops/env.example | 42 ++ .../maria-whatsapp-bridge/ops/install.sh | 158 +++++++ .../ops/maria-dashboard.service | 38 ++ .../ops/maria-rag.service | 48 ++ .../ops/maria-sync.service | 21 + .../ops/maria-sync.timer | 20 + .../ops/maria-whatsapp.service | 47 ++ .../maria-whatsapp-bridge/rag/config.py | 106 +++++ .../maria-whatsapp-bridge/rag/consumer.py | 147 ++++++ .../maria-whatsapp-bridge/rag/indexer.py | 63 +++ .../rag/requirements.txt | 2 + .../maria-whatsapp-bridge/rag/store.py | 53 +++ .../maria-whatsapp-bridge/rag/sync.py | 92 ++++ .../maria-whatsapp-bridge/whatsapp/index.js | 248 +++++++++++ .../whatsapp/package.json | 15 + 20 files changed, 1921 insertions(+) create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/.gitignore create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/api.py create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/index.html create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/login.html create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/env.example create mode 100755 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/install.sh create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-dashboard.service create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-rag.service create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.service create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.timer create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-whatsapp.service create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/config.py create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/consumer.py create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/indexer.py create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/requirements.txt create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/store.py create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/sync.py create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/index.js create mode 100644 proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/package.json diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/.gitignore b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/.gitignore new file mode 100644 index 0000000..727345b --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/.gitignore @@ -0,0 +1,4 @@ +__pycache__/ +*.pyc +node_modules/ +.pytest_cache/ diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md new file mode 100644 index 0000000..0749698 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md @@ -0,0 +1,140 @@ +# Punte WhatsApp + RAG pentru Maria (LXC 171) + +Bot de suport ROA pe WhatsApp: puntea WhatsApp (Baileys) primeste mesajele, +consumer-ul RAG raspunde folosind DOAR informatiile dintr-un depozit de +documente indexat cu embeddings (Ollama), iar un dashboard web controleaza +totul (start/stop/restart, documente, reindexare, sincronizare Google Drive). + +Continua prototipul descris in +`claude-agent/docs/maria-whatsapp-rag-prototype.md` (construit initial in +`/tmp/maria-bridge/`) — aici e mutat in git, ca serviciu persistent. + +Nu confunda cu: +- **Maria pe Flowise** (`vfp_roaauto/COMUN/utile/chatbot/`) — chatbot web separat. +- **Echo / `echo-whatsapp-bridge.service`** (LXC 110 moltbot) — alt bot, alta punte. +- **Punte Discord -> Claude Code** (`discord-bridge/`, acelasi container) — alt + proiect, alt scop (comanda Claude Code de pe Discord, nu suport RAG). + +## Arhitectura + +``` +WhatsApp (self-chat, sau numarul legat) + | + v +whatsapp/index.js (Baileys) -- API HTTP :8099 (/status /send /messages /react /qr) + | + v +rag/consumer.py -- polling la /messages, RAG stateless (FARA memorie intre mesaje) + | vezi docs/maria-whatsapp-rag-prototype.md pentru motiv + v +rag/store.py -- depozit documente (.txt/.md) in ~/.maria-bridge/documents/ +rag/indexer.py -- chunking + embeddings Ollama -> rag_index.json +rag/sync.py -- rclone pull din Google Drive + reindexare conditionata + +dashboard/api.py -- panou web (stdlib, fara dependinte), :18792 + controleaza maria-whatsapp.service + maria-rag.service, + gestioneaza documentele si declanseaza sincronizare/reindexare +``` + +Servicii `systemctl --user` (vezi `ops/`): + +| Unitate | Ce face | +|---|---| +| `maria-whatsapp.service` | Puntea Baileys (Node), port 8099 | +| `maria-rag.service` | Consumer RAG (Python), polling + raspunsuri | +| `maria-sync.service` + `.timer` | Sincronizare Drive + reindexare, la 10 min | +| `maria-dashboard.service` | Panou de control, port 18792 | + +## Instalare + +```bash +cd maria-whatsapp-bridge +./ops/install.sh # creeaza ~/.maria-bridge/, venv, npm install, symlink-uri unit +``` + +Completeaza manual `~/.maria-bridge/env` (copiat din `ops/env.example` la prima +rulare): cel putin `LLM_URL` (backend-ul de chat) si, daca vrei sincronizare +automata cu Drive, `DRIVE_REMOTE`. + +Bridge-ul WhatsApp NU porneste automat la instalare — cere scanarea unui cod QR +(actiune manuala, o singura data): + +```bash +./ops/install.sh --start +# apoi deschide dashboard-ul (tunel SSH catre 127.0.0.1:18792) si scaneaza +# codul QR din cardul "Conectare WhatsApp" +``` + +## Dashboard + +```bash +ssh -L 18792:127.0.0.1:18792 -N claude@10.0.20.171 & +# apoi: http://localhost:18792 +``` + +Autentificare cu `DASHBOARD_TOKEN` din `~/.maria-bridge/env` (generat automat +de `install.sh`). `DASHBOARD_AUTH=off` in env dezactiveaza login-ul — foloseste +DOAR daca panoul ramane strict pe 127.0.0.1/tunel SSH. + +Ce poti face din panou: +- start/stop/restart pentru puntea WhatsApp si consumer-ul RAG +- vezi starea conexiunii WhatsApp si codul QR de asociere (cand nu e conectat) +- listezi, adaugi si stergi documente din depozitul RAG +- reconstruiesti indexul manual, sau declansezi o sincronizare Drive imediata +- citesti ultimele linii din logurile fiecarui serviciu + +## Depozitul de documente + +Fisiere `.txt`/`.md` in `~/.maria-bridge/documents/`. Se pot administra: +1. **manual din dashboard** (adaugare/stergere text, reindexare automata la salvare); +2. **prin sincronizare Google Drive** — vezi mai jos. + +## Sincronizare cu Google Drive + +Sursa: dosarul `D:\GoogleDrive\romfast\document_store` de pe Windows (Google +Drive Desktop). Containerul e headless, deci sincronizarea foloseste +**rclone cu un cont de serviciu** — nu OAuth interactiv in browser. + +Pasi (o singura data): + +1. **Instaleaza rclone** pe container: `sudo apt-get install -y rclone`. +2. **Creeaza un cont de serviciu Google** cu acces la Drive API (Google Cloud + Console -> IAM -> Service Accounts -> Create -> descarca cheia JSON). +3. **Partajeaza folderul** `document_store` din Google Drive cu adresa de email + a contului de serviciu (click dreapta pe folder -> Share), exact cum ai + partaja cu o persoana. Fara acest pas, contul de serviciu nu vede nimic. +4. Pune cheia JSON pe container, ex. `~/.maria-bridge/gdrive-service-account.json` + (0600). +5. Configureaza remote-ul rclone (`rclone config`, fara sesiune interactiva de + browser cu tip `service_account_file`): + ``` + rclone config create gdrive drive \ + scope=drive.readonly \ + service_account_file=/home/claude/.maria-bridge/gdrive-service-account.json + ``` +6. Gaseste ID-ul folderului `document_store` (din URL-ul Drive) si testeaza: + ``` + rclone lsf gdrive: --drive-root-folder-id= + ``` + Sau, mai simplu, foloseste calea prin nume daca folderul e in "My Drive" al + contului care a facut share (rclone urmareste shared-with-me cu + `--drive-shared-with-me` daca e nevoie). +7. Pune tinta gasita in `~/.maria-bridge/env`: + ``` + DRIVE_REMOTE=gdrive:romfast/document_store + ``` +8. Testeaza manual: `systemctl --user start maria-sync.service` apoi + `journalctl --user -u maria-sync -n 50`, sau butonul „sincronizeaza din Drive + acum" din dashboard. + +Dupa configurare, `maria-sync.timer` trage la fiecare 10 minute; reindexarea +ruleaza DOAR daca s-a schimbat efectiv ceva in depozit (amprenta pe nume + +mtime + marime, vezi `rag/sync.py`), ca sa nu reface embeddings degeaba. + +## Context conversational + +`rag/consumer.py` nu retine memorie intre mesaje — fiecare intrebare e o +interogare RAG independenta (system prompt + top-K chunk-uri + intrebare). +Vezi `claude-agent/docs/maria-whatsapp-rag-prototype.md` pentru motiv si +comparatie cu celelalte punti (Discord: context nelimitat + `/new`; Maria pe +Flowise: fereastra fixa de 5 schimburi). diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/api.py b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/api.py new file mode 100644 index 0000000..beb50b8 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/api.py @@ -0,0 +1,419 @@ +#!/usr/bin/env python3 +"""Dashboard de control pentru puntea WhatsApp+RAG a lui Maria (LXC 171). + +Model: `discord-bridge/dashboard/api.py` de pe acelasi container. Diferente +deliberate fata de acela: + + - controleaza DOUA unitati fixe (`maria-whatsapp.service`, `maria-rag.service`), + nu una singura — dar tot dintr-o lista fixa, nu dintr-un nume primit din + request (acelasi motiv: altfel dashboard-ul devine un `systemctl` remote + fara parola). + - expune si un mic depozit de documente (rag/store.py) — listare, adaugare, + stergere — plus un buton de reindexare (ruleaza rag/indexer.py) si starea + de conectare WhatsApp (proxy catre `/status` al puntii, inclusiv QR-ul de + asociere cand nu e inca legata). + - bind implicit pe 127.0.0.1, acelasi model de autentificare prin token/cookie + ca discord-bridge (vezi acolo comentariul complet despre motiv). + +Fara dependinte in afara stdlib, ca sa ruleze fara venv-ul consumer-ului. +""" + +from __future__ import annotations + +import json +import secrets +import subprocess +import sys +import threading +import time +import urllib.error +import urllib.request +from datetime import datetime +from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from urllib.parse import parse_qs, urlparse + +# Radacina proiectului (parintele lui dashboard/) si rag/ trebuie sa fie +# importabile: de acolo vin config.py si store.py. +_DASH = Path(__file__).resolve().parent +_ROOT = _DASH.parent +_RAG = _ROOT / "rag" +for _p in (str(_ROOT), str(_RAG), str(_DASH)): + if _p not in sys.path: + sys.path.insert(0, _p) + +import config # noqa: E402 +import store # noqa: E402 +import sync as sync_mod # noqa: E402 + +# ── constante ─────────────────────────────────────────────────────────── +UNITS = { + "bridge": "maria-whatsapp.service", + "rag": "maria-rag.service", +} +SELF_SERVICE = "maria-dashboard.service" + +COOKIE_NAME = "maria-dashboard" +COOKIE_MAX_AGE = 60 * 60 * 24 * 30 + + +def log_file(which: str) -> Path: + name = {"bridge": "whatsapp.log", "rag": "rag.log", "dashboard": "dashboard.log"}.get(which, "dashboard.log") + return config.LOG_DIR / name + + +def bridge_base_url() -> str: + return f"http://{config.get('BRIDGE_HOST')}:{config.get('BRIDGE_PORT')}" + + +_TOKEN: str | None = None + + +def reset_token_cache() -> None: + global _TOKEN + _TOKEN = None + + +def auth_disabled() -> bool: + return (config.get("DASHBOARD_AUTH") or "").strip().lower() in ("off", "none", "0", "false") + + +def dashboard_token() -> str: + global _TOKEN + if _TOKEN is None: + tok = (config.get("DASHBOARD_TOKEN") or "").strip() + if not tok: + tok = secrets.token_urlsafe(32) + print( + f"[auth] DASHBOARD_TOKEN nesetat in {config.ENV_FILE} — token efemer " + f"pentru acest proces: {tok}", + file=sys.stderr, flush=True, + ) + _TOKEN = tok + return _TOKEN + + +# ── systemd ───────────────────────────────────────────────────────────── +def _sysctl(*args: str, timeout: float = 30.0) -> subprocess.CompletedProcess: + return subprocess.run( + ["systemctl", "--user", *args], + capture_output=True, text=True, timeout=timeout, + ) + + +def _show(unit: str, prop: str) -> str: + try: + return _sysctl("show", "-p", prop, "--value", unit, timeout=5).stdout.strip() + except Exception: + return "" + + +def _uptime_s(unit: str) -> int | None: + mono = _show(unit, "ActiveEnterTimestampMonotonic") + if mono.isdigit() and int(mono) > 0: + return max(0, int(time.monotonic() - int(mono) / 1_000_000)) + ts = _show(unit, "ActiveEnterTimestamp") + if not ts: + return None + try: + parts = ts.split() + stamp = datetime.strptime(f"{parts[1]} {parts[2]}", "%Y-%m-%d %H:%M:%S") + except (ValueError, IndexError): + return None + return max(0, int(time.time() - stamp.timestamp())) + + +def unit_info(unit: str) -> dict: + active = _show(unit, "ActiveState") + pid = _show(unit, "MainPID") + mem = _show(unit, "MemoryCurrent") + restarts = _show(unit, "NRestarts") + return { + "unit": unit, + "active": active == "active", + "state": active or "unknown", + "sub": _show(unit, "SubState"), + "enabled": _show(unit, "UnitFileState"), + "pid": int(pid) if pid.isdigit() and pid != "0" else None, + "memory_bytes": int(mem) if mem.isdigit() else None, + "restarts": int(restarts) if restarts.isdigit() else 0, + "uptime_s": _uptime_s(unit) if active == "active" else None, + } + + +# ── puntea WhatsApp (proxy HTTP catre whatsapp/index.js) ────────────────── +def whatsapp_status() -> dict: + """Interogheaza /status al puntii Node. Esec (proces oprit) = raspuns neutru.""" + try: + req = urllib.request.Request(f"{bridge_base_url()}/status") + with urllib.request.urlopen(req, timeout=3) as resp: + return json.loads(resp.read().decode("utf-8")) + except (urllib.error.URLError, TimeoutError, OSError, ValueError): + return {"connected": False, "phone": None, "qr": None, "reachable": False} + + +# ── depozitul de documente + index ───────────────────────────────────────── +def index_info() -> dict: + try: + raw = config.INDEX_FILE.read_text(encoding="utf-8") + entries = json.loads(raw) + sources = sorted({e.get("source") for e in entries if isinstance(e, dict) and e.get("source")}) + return { + "chunks": len(entries), + "documents_indexed": len(sources), + "mtime": config.INDEX_FILE.stat().st_mtime, + } + except (OSError, ValueError): + return {"chunks": 0, "documents_indexed": 0, "mtime": None} + + +def _venv_python() -> str: + py = config.STATE_DIR / "venv" / "bin" / "python" + return str(py) if py.exists() else sys.executable + + +def run_reindex() -> dict: + """Ruleaza rag/indexer.py cu python-ul din venv (daca exista) si intoarce output-ul.""" + try: + r = subprocess.run( + [_venv_python(), str(_RAG / "indexer.py")], + capture_output=True, text=True, timeout=600, + cwd=str(_RAG), + ) + except subprocess.TimeoutExpired: + return {"ok": False, "error": "indexarea a depasit timpul (600s)"} + ok = r.returncode == 0 + return {"ok": ok, "stdout": r.stdout[-4000:], "stderr": r.stderr[-4000:]} + + +def run_sync(force: bool = False) -> dict: + """Ruleaza rag/sync.py (rclone pull + reindexare conditionata) intr-un proces + separat — acelasi motiv ca la reindex: nu blocam thread-ul dashboard-ului cu + embeddings, iar rclone poate lua cateva secunde bune pe retea.""" + args = [_venv_python(), str(_RAG / "sync.py")] + if force: + args.append("--force") + try: + r = subprocess.run(args, capture_output=True, text=True, timeout=600, cwd=str(_RAG)) + except subprocess.TimeoutExpired: + return {"ok": False, "error": "sincronizarea a depasit timpul (600s)"} + ok = r.returncode == 0 + out = {"ok": ok, "stdout": r.stdout[-4000:], "stderr": r.stderr[-4000:]} + out["state"] = sync_mod.read_sync_state() + return out + + +# ── HTTP ──────────────────────────────────────────────────────────────── +def _parse_cookies(raw: str) -> dict[str, str]: + out: dict[str, str] = {} + for chunk in (raw or "").split(";"): + chunk = chunk.strip() + if "=" in chunk: + k, v = chunk.split("=", 1) + out[k.strip()] = v.strip() + return out + + +class Handler(SimpleHTTPRequestHandler): + server_version = "maria-dashboard" + protocol_version = "HTTP/1.1" + + def __init__(self, *a, **kw): + super().__init__(*a, directory=str(_DASH), **kw) + + def log_message(self, format, *args): # noqa: A002 + sys.stderr.write("%s %s\n" % (self.address_string(), format % args)) + + def send_json(self, payload, status: int = 200, extra_headers: dict | None = None): + body = json.dumps(payload, ensure_ascii=False).encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Content-Length", str(len(body))) + self.send_header("Cache-Control", "no-store") + for k, v in (extra_headers or {}).items(): + self.send_header(k, v) + self.end_headers() + try: + self.wfile.write(body) + except (BrokenPipeError, ConnectionResetError): + pass + + def read_json(self) -> dict: + try: + n = int(self.headers.get("Content-Length") or 0) + raw = self.rfile.read(n).decode("utf-8") if n > 0 else "" + data = json.loads(raw) if raw else {} + return data if isinstance(data, dict) else {} + except (ValueError, OSError, UnicodeDecodeError): + return {} + + def authed(self) -> bool: + if auth_disabled(): + return True + got = _parse_cookies(self.headers.get("Cookie", "")).get(COOKIE_NAME, "") + return bool(got) and secrets.compare_digest(got, dashboard_token()) + + def deny(self): + self.send_json({"error": "neautentificat"}, 401) + + # --- GET ----------------------------------------------------------- + def do_GET(self): + path = urlparse(self.path).path + if path.startswith("/api/"): + if not self.authed(): + return self.deny() + return self.route_get(path) + if path in ("/", "/index.html") and not self.authed(): + self.send_response(302) + self.send_header("Location", "/login.html") + self.send_header("Content-Length", "0") + self.end_headers() + return + if path == "/login.html" and auth_disabled(): + self.send_response(302) + self.send_header("Location", "/") + self.send_header("Content-Length", "0") + self.end_headers() + return + return super().do_GET() + + def route_get(self, path: str): + qs = parse_qs(urlparse(self.path).query) + if path == "/api/status": + return self.send_json({ + "bridge": unit_info(UNITS["bridge"]), + "rag": unit_info(UNITS["rag"]), + "dashboard": unit_info(SELF_SERVICE), + "whatsapp": whatsapp_status(), + "index": index_info(), + "documents": len(store.list_documents()), + "sync": { + **sync_mod.read_sync_state(), + "drive_remote": config.get("DRIVE_REMOTE") or None, + }, + "auth": not auth_disabled(), + "now": time.time(), + }) + if path == "/api/documents": + return self.send_json({"documents": store.list_documents()}) + if path == "/api/documents/read": + name = qs.get("name", [""])[0] + try: + return self.send_json({"name": name, "content": store.read_document(name)}) + except (ValueError, OSError) as exc: + return self.send_json({"error": str(exc)}, 400) + if path == "/api/logs": + try: + n = min(max(int(qs.get("lines", ["200"])[0]), 1), 2000) + except ValueError: + n = 200 + which = qs.get("service", ["bridge"])[0] + target = log_file(which) + if not target.exists(): + return self.send_json({"lines": [f"({target} nu exista)"]}) + r = subprocess.run(["tail", "-n", str(n), str(target)], + capture_output=True, text=True, timeout=15) + return self.send_json({"file": target.name, "lines": r.stdout.splitlines()}) + return self.send_json({"error": "ruta necunoscuta"}, 404) + + # --- POST ---------------------------------------------------------- + def do_POST(self): + path = urlparse(self.path).path + if path == "/api/auth/login": + return self.handle_login() + if path == "/api/auth/logout": + return self.send_json( + {"ok": True}, + extra_headers={"Set-Cookie": f"{COOKIE_NAME}=; HttpOnly; SameSite=Strict; Path=/; Max-Age=0"}, + ) + if not self.authed(): + return self.deny() + if path == "/api/service": + return self.handle_service() + if path == "/api/documents": + return self.handle_document_write() + if path == "/api/documents/delete": + return self.handle_document_delete() + if path == "/api/reindex": + return self.send_json(run_reindex()) + if path == "/api/sync": + data = self.read_json() + return self.send_json(run_sync(force=bool(data.get("force")))) + if path == "/api/restart-self": + self.send_json({"ok": True, "message": "dashboard-ul reporneste in 1s"}) + threading.Thread(target=lambda: (time.sleep(1), __import__("os")._exit(0)), daemon=True).start() + return None + return self.send_json({"error": "ruta necunoscuta"}, 404) + + def handle_login(self): + data = self.read_json() + provided = (data.get("token") or "").strip() + if not provided or not secrets.compare_digest(provided, dashboard_token()): + time.sleep(0.5) + return self.send_json({"error": "token invalid"}, 401) + cookie = (f"{COOKIE_NAME}={dashboard_token()}; HttpOnly; SameSite=Strict; " + f"Path=/; Max-Age={COOKIE_MAX_AGE}") + return self.send_json({"ok": True}, extra_headers={"Set-Cookie": cookie}) + + def handle_service(self): + """start / stop / restart pe una din UNITATILE FIXE. Numele vine dintr-o + cheie (bridge/rag), niciodata dintr-un nume de unit arbitrar.""" + data = self.read_json() + which = str(data.get("service") or "") + action = str(data.get("action") or "") + unit = UNITS.get(which) + if not unit: + return self.send_json({"ok": False, "error": f"serviciu necunoscut: {which}"}, 400) + if action not in ("start", "stop", "restart"): + return self.send_json({"ok": False, "error": f"actiune necunoscuta: {action}"}, 400) + + print(f"[actiune] {action} pe {unit}", file=sys.stderr, flush=True) + try: + r = _sysctl(action, unit) + except subprocess.TimeoutExpired: + return self.send_json({"ok": False, "error": "systemctl a depasit timpul"}, 504) + if r.returncode != 0: + return self.send_json({"ok": False, "error": (r.stderr or r.stdout).strip()}, 500) + time.sleep(1.0) + return self.send_json({"ok": True, "action": action, "service": unit_info(unit)}) + + def handle_document_write(self): + data = self.read_json() + name = str(data.get("name") or "") + content = data.get("content") + if content is None: + return self.send_json({"ok": False, "error": "lipseste 'content'"}, 400) + try: + store.write_document(name, str(content)) + except (ValueError, OSError) as exc: + return self.send_json({"ok": False, "error": str(exc)}, 400) + return self.send_json({"ok": True, "name": name}) + + def handle_document_delete(self): + data = self.read_json() + name = str(data.get("name") or "") + try: + ok = store.delete_document(name) + except (ValueError, OSError) as exc: + return self.send_json({"ok": False, "error": str(exc)}, 400) + return self.send_json({"ok": ok, "name": name}, 200 if ok else 404) + + +def main() -> None: + bind = config.get("DASHBOARD_BIND") or "127.0.0.1" + port = config.get_int("DASHBOARD_PORT", 18792) + dashboard_token() + srv = ThreadingHTTPServer((bind, port), Handler) + srv.daemon_threads = True + print(f"dashboard Maria pe http://{bind}:{port} (unitati: {list(UNITS.values())})", + file=sys.stderr, flush=True) + try: + srv.serve_forever() + except KeyboardInterrupt: + pass + finally: + srv.server_close() + + +if __name__ == "__main__": + main() diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/index.html b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/index.html new file mode 100644 index 0000000..3fa2c38 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/index.html @@ -0,0 +1,211 @@ + + + + + +Maria — WhatsApp + RAG + + + +

Maria — punte WhatsApp + RAG deconectare

+ +
+
+
Punte WhatsApp
+
+ + + +
+
+
Consumer RAG
+
+ + + +
+
+ +

Conectare WhatsApp

+
+
se incarca...
+
+
+ +

Documente indexate (depozit RAG)

+
+
+
+ + +
numemarime
+

Adauga document (.txt / .md)

+ +

+ +

+ +
+
+ +

Loguri

+
+ + +

+
+ + + + diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/login.html b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/login.html new file mode 100644 index 0000000..bbd411c --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/login.html @@ -0,0 +1,47 @@ + + + + + +Maria — autentificare + + + +
+

Maria — dashboard WhatsApp+RAG

+ + +
+
+ + + diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/env.example b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/env.example new file mode 100644 index 0000000..22719a0 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/env.example @@ -0,0 +1,42 @@ +# Configurarea puntii WhatsApp+RAG a lui Maria. +# +# install.sh copiaza acest fisier in ~/.maria-bridge/env cu drepturi 0600. +# Format KEY=value, fara `export`, fara expandare de variabile. + +# --- Puntea WhatsApp (whatsapp/index.js) ----------------------------------- +BRIDGE_HOST=127.0.0.1 +BRIDGE_PORT=8099 +# true = raspunde DOAR in chatul cu tine insuti (self-chat). Pune "false" abia +# dupa ce esti multumit de calitatea raspunsurilor — altfel Maria raspunde la +# oricine iti scrie pe numarul legat. +TEST_MODE_SELF_CHAT_ONLY=true + +# --- Backend LLM + embeddings ----------------------------------------------- +# Modelul de chat folosit pentru raspunsuri (format compatibil OpenAI +# /v1/chat/completions). Implicit presupune un tunel/proxy local catre acelasi +# backend folosit de Echo (LXC 110) — vezi docs/maria-whatsapp-rag-prototype.md. +LLM_URL=http://127.0.0.1:8091 +# Ollama, pentru embeddings (nomic-embed-text). +OLLAMA_URL=http://127.0.0.1:11434 +EMBED_MODEL=nomic-embed-text +TOP_K=3 +MAX_TOKENS=250 +POLL_INTERVAL_S=2 + +# --- Depozit de documente + sincronizare Google Drive ----------------------- +# Tinta rclone pentru dosarul de documente (partajat de pe Windows, +# D:\GoogleDrive\romfast\document_store), format `:`. +# Gol = sincronizare dezactivata; documentele se administreaza doar manual din +# dashboard. Vezi README.md, "Sincronizare cu Google Drive", pentru configurarea +# remote-ului `rclone` (cont de serviciu, fara OAuth interactiv pe un container +# headless). +DRIVE_REMOTE= + +# --- dashboard de control (dashboard/api.py) -------------------------------- +# Token de acces. Generat automat de ops/install.sh daca lipseste. +DASHBOARD_TOKEN= +# `off` scoate complet login-ul. De pus DOAR daca panoul ramane strict pe +# 127.0.0.1 / tailnet (acelasi motiv ca la discord-bridge). +DASHBOARD_AUTH= +DASHBOARD_BIND=127.0.0.1 +DASHBOARD_PORT=18792 diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/install.sh b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/install.sh new file mode 100755 index 0000000..7946dfd --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/install.sh @@ -0,0 +1,158 @@ +#!/usr/bin/env bash +# install.sh — instaleaza puntea WhatsApp+RAG a lui Maria pe LXC 171 (claude-agent). +# +# IDEMPOTENT: se poate rula de cate ori vrei. Nu suprascrie niciodata +# ~/.maria-bridge/env (acolo sta configuratia completata de om). +# +# ./ops/install.sh # instaleaza / actualizeaza, NU porneste serviciile +# ./ops/install.sh --start # in plus porneste bridge+rag (bridge cere QR scanat) +# +# Ce face: +# 1. ~/.maria-bridge/ cu 0700, subdirectoare (logs/, documents/, whatsapp-auth/) +# 2. env 0600 din ops/env.example (doar daca lipseste) +# 3. npm install pentru whatsapp/ (Baileys) +# 4. venv + pip install pentru rag/requirements.txt +# 5. loginctl enable-linger (serviciile de utilizator supravietuiesc logout-ului) +# 6. symlink-uri unit -> ~/.config/systemd/user/ (bridge, rag, dashboard, timer sync) +# 7. systemd-analyze verify + enable +set -uo pipefail + +HERE="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +SRC="$(dirname "$HERE")" # .../maria-whatsapp-bridge +STATE_DIR="$HOME/.maria-bridge" +UNIT_DIR="$HOME/.config/systemd/user" +UNITS=(maria-whatsapp.service maria-rag.service maria-sync.service maria-sync.timer) +DASH_UNIT="maria-dashboard.service" +DO_START=0 +[ "${1:-}" = "--start" ] && DO_START=1 + +info() { printf ' \033[32m*\033[0m %s\n' "$*"; } +warn() { printf ' \033[33m!\033[0m %s\n' "$*"; } +die() { printf ' \033[31mX\033[0m %s\n' "$*" >&2; exit 1; } + +[ "$(id -u)" -eq 0 ] && die "NU rula ca root. Serviciile Maria sunt servicii de utilizator (claude)." + +echo "== Punte WhatsApp+RAG Maria: instalare ==" +echo " sursa: $SRC" + +# --- 1. director de stare -------------------------------------------------- +mkdir -p "$STATE_DIR/logs" "$STATE_DIR/documents" "$STATE_DIR/whatsapp-auth" +chmod 0700 "$STATE_DIR" +info "director de stare: $STATE_DIR (0700)" + +# --- 2. env ---------------------------------------------------------------- +if [ -f "$STATE_DIR/env" ]; then + chmod 0600 "$STATE_DIR/env" + info "env exista deja, nu se atinge" +else + cp "$HERE/env.example" "$STATE_DIR/env" + chmod 0600 "$STATE_DIR/env" + info "env creat din env.example (0600) — completeaza LLM_URL/DRIVE_REMOTE dupa caz" +fi + +# --- 3. dependinte Node (whatsapp/) ----------------------------------------- +if command -v npm >/dev/null 2>&1; then + if [ -d "$SRC/whatsapp/node_modules" ]; then + info "node_modules exista deja in whatsapp/" + else + (cd "$SRC/whatsapp" && npm install --omit=dev) \ + && info "npm install: gata (Baileys si dependinte)" \ + || warn "npm install a esuat (retea?) — reia cu: cd $SRC/whatsapp && npm install" + fi +else + warn "npm lipseste din PATH — instaleaza Node (nvm) inainte de a porni puntea" +fi + +# --- 4. venv + dependinte Python (rag/) ------------------------------------- +if [ ! -x "$STATE_DIR/venv/bin/python" ]; then + python3 -m venv "$STATE_DIR/venv" || die "crearea venv a esuat" + info "venv creat: $STATE_DIR/venv" +else + info "venv exista deja" +fi +"$STATE_DIR/venv/bin/python" -m pip install --quiet --upgrade pip >/dev/null 2>&1 +if "$STATE_DIR/venv/bin/python" -m pip install --quiet -r "$SRC/rag/requirements.txt"; then + info "dependinte Python instalate din rag/requirements.txt" +else + warn "instalarea dependintelor Python a esuat (retea?) — reia cu: $STATE_DIR/venv/bin/pip install -r $SRC/rag/requirements.txt" +fi + +# --- 5. linger --------------------------------------------------------------- +if [ "$(loginctl show-user "$USER" -p Linger --value 2>/dev/null)" = "yes" ]; then + info "linger deja activ pentru $USER" +else + if loginctl enable-linger "$USER" 2>/dev/null; then + info "linger activat pentru $USER" + elif sudo -n loginctl enable-linger "$USER" 2>/dev/null; then + info "linger activat pentru $USER (prin sudo)" + else + warn "nu am putut activa linger; ruleaza manual: sudo loginctl enable-linger $USER" + fi +fi + +# --- 6. unit-uri ------------------------------------------------------------- +mkdir -p "$UNIT_DIR" +for u in "${UNITS[@]}"; do + ln -sfn "$HERE/$u" "$UNIT_DIR/$u" + info "unit legat: $UNIT_DIR/$u -> $HERE/$u" +done +ln -sfn "$SRC/dashboard/$DASH_UNIT" "$UNIT_DIR/$DASH_UNIT" +info "unit legat: $UNIT_DIR/$DASH_UNIT -> $SRC/dashboard/$DASH_UNIT" + +if grep -qE '^DASHBOARD_TOKEN=.+' "$STATE_DIR/env" 2>/dev/null; then + info "DASHBOARD_TOKEN exista deja in env" +else + printf 'DASHBOARD_TOKEN=%s\n' "$(python3 -c 'import secrets;print(secrets.token_urlsafe(24))')" \ + >> "$STATE_DIR/env" + info "DASHBOARD_TOKEN generat si adaugat in $STATE_DIR/env" +fi + +systemctl --user daemon-reload 2>/dev/null || warn "daemon-reload a esuat (sesiune fara systemd de utilizator?)" + +if systemd-analyze verify "$UNIT_DIR/maria-whatsapp.service" 2>&1 | grep -vE 'Unknown key|^$' | grep -q .; then + systemd-analyze verify "$UNIT_DIR/maria-whatsapp.service" 2>&1 | sed 's/^/ /' + warn "systemd-analyze verify a raportat probleme (vezi mai sus)" +else + info "systemd-analyze verify: curat" +fi + +# --- 7. rclone (optional, pentru sincronizare Drive) ------------------------ +if command -v rclone >/dev/null 2>&1; then + info "rclone gasit: $(command -v rclone)" + if ! grep -qE '^DRIVE_REMOTE=.+' "$STATE_DIR/env" 2>/dev/null; then + warn "DRIVE_REMOTE necompletat — sincronizarea cu Google Drive ramane dezactivata" + warn "vezi README.md, sectiunea 'Sincronizare cu Google Drive'" + fi +else + warn "rclone nu e instalat — sincronizarea automata cu Google Drive nu va functiona" + warn "instaleaza cu: sudo apt-get install -y rclone (apoi vezi README.md)" +fi + +# --- 8. enable / start ------------------------------------------------------- +for u in maria-whatsapp.service maria-rag.service; do + systemctl --user enable "$u" >/dev/null 2>&1 \ + && info "$u enabled (porneste la boot)" \ + || warn "enable pentru $u a esuat" +done +systemctl --user enable --now "$DASH_UNIT" >/dev/null 2>&1 \ + && info "dashboard enabled + pornit pe 127.0.0.1:$(grep -E '^DASHBOARD_PORT=' "$STATE_DIR/env" | cut -d= -f2- || echo 18792)" \ + || warn "enable pentru $DASH_UNIT a esuat" +systemctl --user enable --now maria-sync.timer >/dev/null 2>&1 \ + && info "timer de sincronizare Drive activat (la fiecare 10 min)" \ + || warn "enable pentru maria-sync.timer a esuat" + +if [ "$DO_START" -eq 1 ]; then + systemctl --user restart maria-whatsapp.service && info "punte WhatsApp pornita" + echo + echo " Deschide dashboard-ul (127.0.0.1:\$DASHBOARD_PORT, tunel SSH) si scaneaza" + echo " codul QR din cardul 'Conectare WhatsApp' inainte sa pornesti maria-rag." + systemctl --user restart maria-rag.service && info "consumer RAG pornit" + systemctl --user --no-pager status maria-whatsapp.service | sed 's/^/ /' +else + echo + echo " Serviciile bridge/rag NU au fost pornite (intentionat: bridge-ul cere" + echo " scanarea unui cod QR de asociere WhatsApp, actiune manuala)." + echo " Dupa ce completezi $STATE_DIR/env: $0 --start" +fi + +echo "== gata ==" diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-dashboard.service b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-dashboard.service new file mode 100644 index 0000000..c20de86 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-dashboard.service @@ -0,0 +1,38 @@ +# Unit systemd de UTILIZATOR pentru dashboard-ul de control al puntii Maria. +# +# Se instaleaza in ~/.config/systemd/user/maria-dashboard.service (vezi install.sh). +# Unit SEPARAT de bridge/rag, ca o repornire a lor sa nu ia si panoul din care +# ai apasat butonul (acelasi motiv ca la claude-discord-dashboard.service). +# +# systemctl --user daemon-reload +# systemctl --user enable --now maria-dashboard +# journalctl --user -u maria-dashboard -f + +[Unit] +Description=Dashboard de control pentru puntea WhatsApp+RAG a lui Maria (LXC 171) +Documentation=file:///workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md +After=network.target + +[Service] +Type=simple +WorkingDirectory=%h/.maria-bridge +EnvironmentFile=%h/.maria-bridge/env +Environment=MARIA_BRIDGE_DIR=%h/.maria-bridge +Environment=PYTHONUNBUFFERED=1 + +# Stdlib only, dar are nevoie sa dea `systemctl --user` — acelasi utilizator ca +# bridge/rag, deci merge fara sudo. +ExecStart=/usr/bin/python3 /workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/dashboard/api.py + +Restart=always +RestartSec=2 + +MemoryHigh=192M +MemoryMax=384M + +StandardOutput=append:%h/.maria-bridge/logs/dashboard.log +StandardError=append:%h/.maria-bridge/logs/dashboard.log +SyslogIdentifier=maria-dashboard + +[Install] +WantedBy=default.target diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-rag.service b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-rag.service new file mode 100644 index 0000000..c34d51d --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-rag.service @@ -0,0 +1,48 @@ +# Unit systemd de UTILIZATOR pentru consumer-ul RAG al lui Maria (rag/consumer.py). +# +# Se instaleaza in ~/.config/systemd/user/maria-rag.service (vezi install.sh). +# +# systemctl --user daemon-reload +# systemctl --user enable --now maria-rag +# journalctl --user -u maria-rag -f + +[Unit] +Description=Consumer RAG pentru Maria — raspunde la mesajele din puntea WhatsApp +Documentation=file:///workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md +After=network-online.target maria-whatsapp.service +Wants=network-online.target +# Nu e o dependinta stricta (poate porni si daca puntea WhatsApp inca nu e sus — +# consumer.py doar face polling si tolereaza erori de retea), dar ordinea are sens. +Wants=maria-whatsapp.service + +StartLimitIntervalSec=300 +StartLimitBurst=5 + +[Service] +Type=simple +WorkingDirectory=/workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag +EnvironmentFile=%h/.maria-bridge/env +Environment=MARIA_BRIDGE_DIR=%h/.maria-bridge +Environment=PYTHONUNBUFFERED=1 + +ExecStartPre=/bin/mkdir -p %h/.maria-bridge/logs %h/.maria-bridge/documents +ExecStart=%h/.maria-bridge/venv/bin/python consumer.py + +KillMode=control-group +KillSignal=SIGTERM +TimeoutStopSec=15 + +Restart=on-failure +RestartSec=5 +RestartSteps=5 +RestartMaxDelaySec=120 + +MemoryHigh=256M +MemoryMax=512M + +StandardOutput=append:%h/.maria-bridge/logs/rag.log +StandardError=append:%h/.maria-bridge/logs/rag.log +SyslogIdentifier=maria-rag + +[Install] +WantedBy=default.target diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.service b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.service new file mode 100644 index 0000000..1b30aae --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.service @@ -0,0 +1,21 @@ +# Unit systemd de UTILIZATOR: o singura trecere de sincronizare Drive + reindexare +# conditionata (rag/sync.py). Pornit periodic de maria-sync.timer, sau manual: +# +# systemctl --user start maria-sync.service +# journalctl --user -u maria-sync -f + +[Unit] +Description=Sincronizare depozit documente Maria (Drive) + reindexare RAG +Documentation=file:///workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md + +[Service] +Type=oneshot +WorkingDirectory=/workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag +EnvironmentFile=%h/.maria-bridge/env +Environment=MARIA_BRIDGE_DIR=%h/.maria-bridge +Environment=PYTHONUNBUFFERED=1 +ExecStart=%h/.maria-bridge/venv/bin/python sync.py + +StandardOutput=append:%h/.maria-bridge/logs/sync.log +StandardError=append:%h/.maria-bridge/logs/sync.log +SyslogIdentifier=maria-sync diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.timer b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.timer new file mode 100644 index 0000000..ae87a85 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-sync.timer @@ -0,0 +1,20 @@ +# Timer de UTILIZATOR: declanseaza maria-sync.service la fiecare 10 minute. +# +# Interval fix aici (nu citit din env): un timer systemd nu poate parametriza +# OnUnitActiveSec dintr-un fisier extern usor — de editat direct daca 10 min nu +# se potriveste. +# +# systemctl --user daemon-reload +# systemctl --user enable --now maria-sync.timer +# systemctl --user list-timers maria-sync.timer + +[Unit] +Description=Sincronizare periodica a depozitului de documente Maria (Drive) + +[Timer] +OnBootSec=2min +OnUnitActiveSec=10min +Persistent=true + +[Install] +WantedBy=timers.target diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-whatsapp.service b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-whatsapp.service new file mode 100644 index 0000000..56b0fdb --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/ops/maria-whatsapp.service @@ -0,0 +1,47 @@ +# Unit systemd de UTILIZATOR pentru puntea WhatsApp (Baileys) a lui Maria. +# +# Se instaleaza in ~/.config/systemd/user/maria-whatsapp.service (vezi install.sh). +# +# systemctl --user daemon-reload +# systemctl --user enable --now maria-whatsapp +# journalctl --user -u maria-whatsapp -f + +[Unit] +Description=Punte WhatsApp (Baileys) pentru Maria (LXC 171 claude-agent) +Documentation=file:///workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/README.md +After=network-online.target +Wants=network-online.target + +StartLimitIntervalSec=300 +StartLimitBurst=5 + +[Service] +Type=simple +WorkingDirectory=/workspace/romfastsql/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp +EnvironmentFile=%h/.maria-bridge/env +Environment=MARIA_BRIDGE_DIR=%h/.maria-bridge +Environment=NODE_ENV=production +# nvm — acelasi PATH ca discord-bridge, pe acelasi container. +Environment=PATH=%h/bin:%h/.nvm/versions/node/v20.19.6/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin + +ExecStartPre=/bin/mkdir -p %h/.maria-bridge/logs %h/.maria-bridge/whatsapp-auth +ExecStart=%h/.nvm/versions/node/v20.19.6/bin/node index.js + +KillMode=control-group +KillSignal=SIGTERM +TimeoutStopSec=15 + +Restart=on-failure +RestartSec=5 +RestartSteps=5 +RestartMaxDelaySec=120 + +MemoryHigh=384M +MemoryMax=768M + +StandardOutput=append:%h/.maria-bridge/logs/whatsapp.log +StandardError=append:%h/.maria-bridge/logs/whatsapp.log +SyslogIdentifier=maria-whatsapp + +[Install] +WantedBy=default.target diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/config.py b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/config.py new file mode 100644 index 0000000..b765f49 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/config.py @@ -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() diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/consumer.py b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/consumer.py new file mode 100644 index 0000000..f7d6f9d --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/consumer.py @@ -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() diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/indexer.py b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/indexer.py new file mode 100644 index 0000000..3370e98 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/indexer.py @@ -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, + ) diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/requirements.txt b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/requirements.txt new file mode 100644 index 0000000..0912bc7 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/requirements.txt @@ -0,0 +1,2 @@ +# Consumer + indexer RAG pentru Maria. Restul e stdlib. +requests==2.32.3 diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/store.py b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/store.py new file mode 100644 index 0000000..e46f5b7 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/store.py @@ -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 diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/sync.py b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/sync.py new file mode 100644 index 0000000..615c533 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/rag/sync.py @@ -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 -> 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) diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/index.js b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/index.js new file mode 100644 index 0000000..72883f6 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/index.js @@ -0,0 +1,248 @@ +// Punte WhatsApp (Baileys) pentru Maria — implicit doar self-chat (vezi +// TEST_MODE_SELF_CHAT_ONLY). Sesiunea WhatsApp (auth/) NU intra in git — +// contine chei de sesiune, e in ~/.maria-bridge/whatsapp-auth/ (vezi config.py). + +const express = require('express'); +const pino = require('pino'); +const QRCode = require('qrcode'); +const path = require('path'); + +let makeWASocket, useMultiFileAuthState, DisconnectReason, fetchLatestBaileysVersion; + +const PORT = parseInt(process.env.BRIDGE_PORT || '8099', 10); +const HOST = process.env.BRIDGE_HOST || '127.0.0.1'; +const AUTH_DIR = process.env.MARIA_BRIDGE_DIR + ? path.join(process.env.MARIA_BRIDGE_DIR, 'whatsapp-auth') + : path.join(__dirname, 'auth'); +const MAX_RECONNECT_ATTEMPTS = 5; +const TEST_MODE_SELF_CHAT_ONLY = (process.env.TEST_MODE_SELF_CHAT_ONLY || 'true') !== 'false'; + +const logger = pino({ level: 'warn' }); + +let sock = null; +let connected = false; +let phoneNumber = null; +let ownJid = null; +let currentQR = null; +let currentPairingCode = null; +let reconnectAttempts = 0; +let messageQueue = []; +let shuttingDown = false; + +async function startConnection() { + if (!makeWASocket) { + const baileys = await import('@whiskeysockets/baileys'); + makeWASocket = baileys.default; + useMultiFileAuthState = baileys.useMultiFileAuthState; + DisconnectReason = baileys.DisconnectReason; + fetchLatestBaileysVersion = baileys.fetchLatestBaileysVersion; + } + + const { state, saveCreds } = await useMultiFileAuthState(AUTH_DIR); + const { version } = await fetchLatestBaileysVersion(); + + sock = makeWASocket({ + version, + auth: state, + logger, + printQRInTerminal: false, + defaultQueryTimeoutMs: 60000, + }); + + sock.ev.on('creds.update', saveCreds); + + sock.ev.on('connection.update', async (update) => { + const { connection, lastDisconnect, qr } = update; + + if (qr) { + try { + currentQR = await QRCode.toDataURL(qr); + console.log('[whatsapp] QR code generated — scan with WhatsApp to link'); + } catch (err) { + console.error('[whatsapp] Failed to generate QR code:', err.message); + } + } + + if (connection === 'open') { + connected = true; + currentQR = null; + reconnectAttempts = 0; + phoneNumber = sock.user?.id?.split(':')[0] || sock.user?.id?.split('@')[0] || null; + ownJid = phoneNumber ? `${phoneNumber}@s.whatsapp.net` : null; + console.log(`[whatsapp] Connected as ${phoneNumber} (self-chat-only mode: ${TEST_MODE_SELF_CHAT_ONLY})`); + } + + if (connection === 'close') { + connected = false; + phoneNumber = null; + const statusCode = lastDisconnect?.error?.output?.statusCode; + const shouldReconnect = statusCode !== DisconnectReason.loggedOut; + + console.log(`[whatsapp] Disconnected (status: ${statusCode})`); + + if (shouldReconnect && !shuttingDown) { + if (reconnectAttempts < MAX_RECONNECT_ATTEMPTS) { + reconnectAttempts++; + const delay = Math.min(1000 * Math.pow(2, reconnectAttempts), 30000); + console.log(`[whatsapp] Reconnecting in ${delay}ms (attempt ${reconnectAttempts}/${MAX_RECONNECT_ATTEMPTS})`); + setTimeout(startConnection, delay); + } else { + console.error(`[whatsapp] Max reconnect attempts reached (${MAX_RECONNECT_ATTEMPTS})`); + } + } else if (statusCode === DisconnectReason.loggedOut) { + console.log('[whatsapp] Logged out — delete auth dir and restart to re-link'); + } + } + }); + + sock.ev.on('messages.upsert', ({ messages, type }) => { + if (type !== 'notify') return; + + for (const msg of messages) { + if (msg.key.remoteJid === 'status@broadcast') continue; + const isGroup = msg.key.remoteJid.endsWith('@g.us'); + const isSelfChat = ownJid && msg.key.remoteJid === ownJid; + if (TEST_MODE_SELF_CHAT_ONLY && !isSelfChat) continue; + if (msg.key.fromMe && !isGroup && !isSelfChat) continue; + const text = msg.message?.conversation || msg.message?.extendedTextMessage?.text; + if (!text) continue; + + messageQueue.push({ + from: msg.key.remoteJid, + participant: msg.key.participant || null, + pushName: msg.pushName || null, + text, + timestamp: msg.messageTimestamp, + id: msg.key.id, + isGroup, + fromMe: msg.key.fromMe || false, + }); + + console.log(`[whatsapp] Message from ${msg.pushName || 'unknown'} in ${msg.key.remoteJid}: ${text.substring(0, 80)}`); + } + }); +} + +const app = express(); +app.use(express.json({ limit: '50mb' })); + +app.get('/status', (_req, res) => { + res.json({ + connected, + phone: phoneNumber, + selfChatOnly: TEST_MODE_SELF_CHAT_ONLY, + qr: connected ? null : currentQR, + }); +}); + +app.post('/pair', async (req, res) => { + if (connected) { + return res.json({ error: 'already connected' }); + } + const { phone } = req.body || {}; + if (!phone) { + return res.status(400).json({ error: 'missing "phone" in body' }); + } + if (!sock) { + return res.status(503).json({ error: 'socket not ready yet, try again in a few seconds' }); + } + try { + const code = await sock.requestPairingCode(phone.replace(/\D/g, '')); + currentPairingCode = code; + console.log(`[whatsapp] Pairing code for ${phone}: ${code}`); + res.json({ ok: true, code }); + } catch (err) { + console.error('[whatsapp] Pairing code error:', err.message); + res.status(500).json({ error: err.message }); + } +}); + +app.get('/pair-code', (_req, res) => { + if (connected) return res.json({ error: 'already connected' }); + if (!currentPairingCode) return res.json({ error: 'no pairing code yet — POST /pair first' }); + res.json({ code: currentPairingCode }); +}); + +app.get('/qr', (_req, res) => { + if (connected) { + return res.json({ error: 'already connected' }); + } + if (!currentQR) { + return res.json({ error: 'no QR code available yet' }); + } + const html = ` +WhatsApp QR + + +QR Code

Scan with WhatsApp → Linked Devices

`; + res.type('html').send(html); +}); + +app.post('/send', async (req, res) => { + const { to, text } = req.body || {}; + + if (!to || !text) { + return res.status(400).json({ ok: false, error: 'missing "to" or "text" in body' }); + } + if (!connected || !sock) { + return res.status(503).json({ ok: false, error: 'not connected to WhatsApp' }); + } + + try { + const result = await sock.sendMessage(to, { text }); + res.json({ ok: true, id: result.key.id }); + } catch (err) { + console.error('[whatsapp] Send failed:', err.message); + res.status(500).json({ ok: false, error: err.message }); + } +}); + +app.post('/react', async (req, res) => { + const { to, id, emoji, fromMe, participant } = req.body || {}; + + if (!to || !id || emoji == null) { + return res.status(400).json({ ok: false, error: 'missing "to", "id", or "emoji" in body' }); + } + if (!connected || !sock) { + return res.status(503).json({ ok: false, error: 'not connected to WhatsApp' }); + } + + try { + const key = { remoteJid: to, id, fromMe: fromMe || false }; + if (participant) key.participant = participant; + await sock.sendMessage(to, { react: { text: emoji, key } }); + res.json({ ok: true }); + } catch (err) { + console.error('[whatsapp] React failed:', err.message); + res.status(500).json({ ok: false, error: err.message }); + } +}); + +app.get('/messages', (_req, res) => { + const messages = messageQueue.splice(0); + res.json({ messages }); +}); + +const server = app.listen(PORT, HOST, () => { + console.log(`[whatsapp] Bridge API listening on http://${HOST}:${PORT}`); + startConnection().catch((err) => { + console.error('[whatsapp] Failed to start connection:', err.message); + }); +}); + +function shutdown(signal) { + console.log(`[whatsapp] Received ${signal}, shutting down...`); + shuttingDown = true; + if (sock) { + sock.end(undefined); + } + server.close(() => { + console.log('[whatsapp] HTTP server closed'); + process.exit(0); + }); + setTimeout(() => process.exit(1), 5000); +} + +process.on('SIGTERM', () => shutdown('SIGTERM')); +process.on('SIGINT', () => shutdown('SIGINT')); diff --git a/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/package.json b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/package.json new file mode 100644 index 0000000..283f6f4 --- /dev/null +++ b/proxmox/lxc171-claude-agent/maria-whatsapp-bridge/whatsapp/package.json @@ -0,0 +1,15 @@ +{ + "name": "maria-whatsapp-bridge", + "version": "1.0.0", + "description": "WhatsApp bridge for Maria RAG support bot, using Baileys", + "main": "index.js", + "scripts": { + "start": "node index.js" + }, + "dependencies": { + "@whiskeysockets/baileys": "^6.7.16", + "express": "^4.21.0", + "pino": "^9.6.0", + "qrcode": "^1.5.4" + } +}