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>
420 lines
16 KiB
Python
420 lines
16 KiB
Python
#!/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()
|