Un client poate avea mai multi angajati care scriu pe WhatsApp de pe numere diferite; toti gasesc acum acelasi profil (rag/client.py, ~/.maria-bridge/clients.json). Profilul intra in prompt sub o sectiune DESPRE CLIENT, la fiecare mesaj, stateless ca restul Mariei — dar nu atinge cautarea/scorul din rag/rank.py, ca sa nu strice pragurile deja calibrate. Administrare din dashboard (sectiune noua "Clienti cunoscuti"), nu din chat, ca sa ramana pe canalul autentificat. 117/117 teste maria-whatsapp-bridge + teste noi pentru rutele de dashboard. Co-Authored-By: Claude Agent <noreply@anthropic.com>
1013 lines
40 KiB
Python
1013 lines
40 KiB
Python
#!/usr/bin/env python3
|
|
"""Dashboard de control pentru puntea Discord -> Claude Code (LXC 171).
|
|
|
|
Model: `echo-core/dashboard` de pe LXC 110 (server stdlib + handler-ul `eco.py`).
|
|
Diferentele deliberate:
|
|
|
|
- **o singura unitate controlata**: `claude-discord.service`. Nu exista endpoint
|
|
care sa primeasca un nume de unit din exterior — altfel dashboard-ul ar deveni
|
|
un `systemctl` remote fara parola.
|
|
- **bind pe 127.0.0.1 implicit**: butonul "restart" opreste un agent care ruleaza
|
|
cu `bypassPermissions` si chei SSH catre tot clusterul. Accesul se face prin
|
|
tunel SSH (vezi README), nu expus in LAN.
|
|
- **restart protejat de tururi in zbor**: `stop`/`restart` intorc 409 daca exista
|
|
fire cu tur in desfasurare, pana cand se cere explicit `force`.
|
|
|
|
Fara dependinte in afara stdlib: ruleaza cu acelasi python ca botul, dar nu are
|
|
nevoie de venv-ul lui.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import os
|
|
import re
|
|
import secrets
|
|
import shutil
|
|
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 punti (parintele lui dashboard/) trebuie sa fie importabila: de acolo
|
|
# vin config.py, cleanup.py si security/approvals.py.
|
|
_DASH = Path(__file__).resolve().parent
|
|
_BRIDGE = _DASH.parent
|
|
for _p in (str(_BRIDGE), str(_DASH)):
|
|
if _p not in sys.path:
|
|
sys.path.insert(0, _p)
|
|
|
|
import config # noqa: E402
|
|
from limits import parse_cap # noqa: E402
|
|
|
|
# ── constante ───────────────────────────────────────────────────────────
|
|
SERVICE = "claude-discord.service"
|
|
SELF_SERVICE = "claude-discord-dashboard.service"
|
|
|
|
# ── Maria (WhatsApp+RAG, proiect sibling in acelasi repo) ────────────────
|
|
# Panou unic la cererea operatorului: acest dashboard controleaza si serviciile
|
|
# Mariei, nu doar puntea Discord. NU importam rag/config.py de acolo (ar coliza
|
|
# pe numele de modul "config" cu discord-bridge/config.py, deja importat mai
|
|
# sus) — citim direct caile si folosim subprocess/HTTP, exact ca la SERVICE.
|
|
MARIA_UNITS = {"bridge": "maria-whatsapp.service", "rag": "maria-rag.service"}
|
|
# Oglinda lui rag/store.py: DOC_EXTENSIONS si ordinea de preferinta (.xml cel mai bogat).
|
|
MARIA_DOC_EXTENSIONS = (".txt", ".md", ".xml")
|
|
_MARIA_FORMAT_RANK = {".xml": 0, ".md": 1, ".txt": 2}
|
|
MARIA_STATE_DIR = Path(os.environ.get("MARIA_BRIDGE_DIR") or (Path.home() / ".maria-bridge"))
|
|
MARIA_DOCS_DIR = MARIA_STATE_DIR / "documents"
|
|
# Documentele care nu vin din Drive (rclone sync sterge din MARIA_DOCS_DIR tot ce
|
|
# nu exista in Drive). Oglinda lui rag/config.py:DOCS_LOCAL_DIR.
|
|
MARIA_DOCS_LOCAL_DIR = MARIA_STATE_DIR / "documents-local"
|
|
MARIA_INDEX_FILE = MARIA_STATE_DIR / "rag_index.json"
|
|
MARIA_ENV_FILE = MARIA_STATE_DIR / "env"
|
|
MARIA_LOG_DIR = MARIA_STATE_DIR / "logs"
|
|
MARIA_VENV_PY = MARIA_STATE_DIR / "venv" / "bin" / "python"
|
|
MARIA_RAG_DIR = _BRIDGE.parent / "maria-whatsapp-bridge" / "rag"
|
|
_MARIA_SAFE_NAME = re.compile(r"^[A-Za-z0-9._-]{1,200}$")
|
|
MARIA_CLIENTS_FILE = MARIA_STATE_DIR / "clients.json"
|
|
_MARIA_CLIENT_ID = re.compile(r"^[a-z0-9_-]{1,40}$")
|
|
|
|
|
|
def maria_env() -> dict:
|
|
try:
|
|
return config.parse_env(MARIA_ENV_FILE.read_text(encoding="utf-8"))
|
|
except OSError:
|
|
return {}
|
|
|
|
|
|
def maria_get(key: str, default: str = "") -> str:
|
|
return maria_env().get(key, default)
|
|
|
|
|
|
def maria_bridge_url() -> str:
|
|
return f"http://{maria_get('BRIDGE_HOST', '127.0.0.1')}:{maria_get('BRIDGE_PORT', '8099')}"
|
|
|
|
|
|
def maria_log(which: str) -> Path:
|
|
name = {"bridge": "whatsapp.log", "rag": "rag.log"}.get(which, "rag.log")
|
|
return MARIA_LOG_DIR / name
|
|
|
|
|
|
def maria_whatsapp_status() -> dict:
|
|
try:
|
|
req = urllib.request.Request(f"{maria_bridge_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}
|
|
|
|
|
|
def maria_request_pairing_code(phone: str) -> tuple[dict, int]:
|
|
"""Cere puntii un cod de asociere de 8 caractere pentru `phone`.
|
|
|
|
Dashboard-ul nu decide nimic despre numar — puntea valideaza si raspunde; noi
|
|
doar transmitem, ca sa nu existe doua reguli de validare care se pot desincroniza.
|
|
"""
|
|
body = json.dumps({"phone": phone}).encode("utf-8")
|
|
req = urllib.request.Request(
|
|
f"{maria_bridge_url()}/pair", data=body,
|
|
headers={"Content-Type": "application/json"}, method="POST",
|
|
)
|
|
try:
|
|
with urllib.request.urlopen(req, timeout=30) as resp:
|
|
return json.loads(resp.read().decode("utf-8")), resp.status
|
|
except urllib.error.HTTPError as exc:
|
|
try:
|
|
return json.loads(exc.read().decode("utf-8")), exc.code
|
|
except ValueError:
|
|
return {"ok": False, "error": f"puntea a raspuns {exc.code}"}, exc.code
|
|
except (urllib.error.URLError, TimeoutError, OSError, ValueError) as exc:
|
|
return {"ok": False, "error": f"puntea WhatsApp nu raspunde: {exc}"}, 503
|
|
|
|
|
|
def maria_index_info() -> dict:
|
|
try:
|
|
entries = json.loads(MARIA_INDEX_FILE.read_text(encoding="utf-8"))
|
|
sources = {e.get("source") for e in entries if isinstance(e, dict) and e.get("source")}
|
|
return {"chunks": len(entries), "documents_indexed": len(sources),
|
|
"mtime": MARIA_INDEX_FILE.stat().st_mtime}
|
|
except (OSError, ValueError):
|
|
return {"chunks": 0, "documents_indexed": 0, "mtime": None}
|
|
|
|
|
|
def maria_sync_state() -> dict:
|
|
try:
|
|
return json.loads((MARIA_STATE_DIR / ".sync_state.json").read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
return {}
|
|
|
|
|
|
def maria_documents() -> list[dict]:
|
|
"""Oglinda lui `maria-whatsapp-bridge/rag/store.py:list_documents` — tine-le la fel.
|
|
|
|
Nu importam modulul acela: si el, si `discord-bridge/config.py` s-ar numi `config`.
|
|
"""
|
|
MARIA_DOCS_DIR.mkdir(parents=True, exist_ok=True)
|
|
vazute: dict[str, Path] = {}
|
|
for d in (MARIA_DOCS_LOCAL_DIR, MARIA_DOCS_DIR): # Drive ultimul = castiga la acelasi nume
|
|
if not d.is_dir():
|
|
continue
|
|
for f in sorted(d.glob("*")):
|
|
if f.is_file() and f.suffix in MARIA_DOC_EXTENSIONS:
|
|
vazute[f.name] = f
|
|
files = [vazute[n] for n in sorted(vazute)]
|
|
winners: dict[str, str] = {}
|
|
for f in files:
|
|
best = winners.get(f.stem)
|
|
if best is None or _MARIA_FORMAT_RANK[f.suffix] < _MARIA_FORMAT_RANK[Path(best).suffix]:
|
|
winners[f.stem] = f.name
|
|
|
|
out = []
|
|
for f in files:
|
|
st = f.stat()
|
|
winner = winners[f.stem]
|
|
out.append({
|
|
"name": f.name,
|
|
"size": st.st_size,
|
|
"mtime": st.st_mtime,
|
|
"local": f.parent == MARIA_DOCS_LOCAL_DIR,
|
|
"shadowed_by": None if winner == f.name else winner,
|
|
})
|
|
return out
|
|
|
|
|
|
def maria_escalations(limit: int = 12) -> list[dict]:
|
|
"""Intrebarile pe care Maria le-a trimis la suport, cele mai noi primele.
|
|
|
|
Sursa: fisierele scrise de `rag/consumer.py:escalate`. Se citesc de pe disc la
|
|
fiecare cerere — sunt putine si e mai bine decat inca o stare de tinut sincron.
|
|
"""
|
|
d = MARIA_STATE_DIR / "escalations"
|
|
if not d.is_dir():
|
|
return []
|
|
out = []
|
|
for f in sorted(d.glob("*.json"), reverse=True)[:limit]:
|
|
try:
|
|
rec = json.loads(f.read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
continue
|
|
out.append({
|
|
"ref": rec.get("ref"),
|
|
"at": rec.get("at"),
|
|
"from": rec.get("push_name") or rec.get("from"),
|
|
"reason": rec.get("reason"),
|
|
"best_cosine": rec.get("best_cosine"),
|
|
"had_image": rec.get("had_image"),
|
|
"notified": rec.get("notified"),
|
|
"notify_error": rec.get("notify_error"),
|
|
# textul citit din captura, daca a fost una: exact ce a "vazut" Maria
|
|
"text": (rec.get("ocr_text") or rec.get("text") or "")[:600],
|
|
})
|
|
return out
|
|
|
|
|
|
def _maria_validate_name(name: str) -> str:
|
|
if not name or not _MARIA_SAFE_NAME.match(name) or ".." in name or "/" in name:
|
|
raise ValueError(f"nume de document invalid: {name!r}")
|
|
if not name.endswith(MARIA_DOC_EXTENSIONS):
|
|
raise ValueError("doar fisiere " + ", ".join(MARIA_DOC_EXTENSIONS))
|
|
return name
|
|
|
|
|
|
def maria_write_document(name: str, content: str) -> None:
|
|
"""Scrie in documents-local/, nu in oglinda Drive.
|
|
|
|
Pana acum ajungea in `documents/` si disparea la urmatorul `rclone sync` —
|
|
documentul adaugat din dashboard traia pana la 10 minute.
|
|
"""
|
|
name = _maria_validate_name(name)
|
|
MARIA_DOCS_LOCAL_DIR.mkdir(parents=True, exist_ok=True)
|
|
(MARIA_DOCS_LOCAL_DIR / name).write_text(content, encoding="utf-8")
|
|
|
|
|
|
def maria_delete_document(name: str) -> bool:
|
|
"""Sterge documentul din oricare din cele doua directoare.
|
|
|
|
Cel din Drive revine la urmatoarea sincronizare — sterge-l din Drive daca vrei
|
|
sa dispara definitiv.
|
|
"""
|
|
name = _maria_validate_name(name)
|
|
for d in (MARIA_DOCS_DIR, MARIA_DOCS_LOCAL_DIR):
|
|
path = d / name
|
|
if path.exists():
|
|
path.unlink()
|
|
return True
|
|
return False
|
|
|
|
|
|
def maria_clients() -> list[dict]:
|
|
"""Oglinda lui `maria-whatsapp-bridge/rag/client.py:toti` — tine-o la fel.
|
|
|
|
Nu importam modulul acela (colizeaza pe "config", ca la documente).
|
|
"""
|
|
try:
|
|
date = json.loads(MARIA_CLIENTS_FILE.read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
date = {}
|
|
return [{"id": cid, **c} for cid, c in sorted(date.items())]
|
|
|
|
|
|
def _maria_save_clients(date: dict) -> None:
|
|
MARIA_STATE_DIR.mkdir(parents=True, exist_ok=True)
|
|
tmp = MARIA_CLIENTS_FILE.with_suffix(".tmp")
|
|
tmp.write_text(json.dumps(date, ensure_ascii=False, indent=2), encoding="utf-8")
|
|
os.replace(tmp, MARIA_CLIENTS_FILE)
|
|
|
|
|
|
def maria_save_client(cid: str, nume: str, profil: str, numere: list[str]) -> None:
|
|
if not cid or not _MARIA_CLIENT_ID.match(cid):
|
|
raise ValueError(f"id de client invalid: {cid!r}")
|
|
normalizate = [n for n in (re.sub(r"\D", "", x) for x in numere) if n]
|
|
try:
|
|
date = json.loads(MARIA_CLIENTS_FILE.read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
date = {}
|
|
for alt_id, alt in date.items():
|
|
if alt_id == cid:
|
|
continue
|
|
furate = set(normalizate) & set(alt.get("numere", []))
|
|
if furate:
|
|
raise ValueError(
|
|
f"numarul {sorted(furate)[0]} e deja atribuit clientului {alt.get('nume', alt_id)!r}")
|
|
date[cid] = {"nume": nume, "profil": profil, "numere": normalizate}
|
|
_maria_save_clients(date)
|
|
|
|
|
|
def maria_delete_client(cid: str) -> bool:
|
|
try:
|
|
date = json.loads(MARIA_CLIENTS_FILE.read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
date = {}
|
|
if cid not in date:
|
|
return False
|
|
del date[cid]
|
|
_maria_save_clients(date)
|
|
return True
|
|
|
|
|
|
def _maria_run(script: str, timeout: float = 600.0, extra_args: list[str] | None = None) -> dict:
|
|
interpreter = str(MARIA_VENV_PY) if MARIA_VENV_PY.exists() else sys.executable
|
|
args = [interpreter, str(MARIA_RAG_DIR / script), *(extra_args or [])]
|
|
try:
|
|
r = subprocess.run(args, capture_output=True, text=True, timeout=timeout, cwd=str(MARIA_RAG_DIR))
|
|
except subprocess.TimeoutExpired:
|
|
return {"ok": False, "error": f"{script} a depasit timpul ({timeout:.0f}s)"}
|
|
return {"ok": r.returncode == 0, "stdout": r.stdout[-4000:], "stderr": r.stderr[-4000:]}
|
|
|
|
|
|
def maria_reindex() -> dict:
|
|
return _maria_run("indexer.py")
|
|
|
|
|
|
def maria_sync(force: bool = False) -> dict:
|
|
out = _maria_run("sync.py", extra_args=(["--force"] if force else []))
|
|
out["state"] = maria_sync_state()
|
|
return out
|
|
|
|
|
|
# Caile se recalculeaza la fiecare apel, nu se ingheata la import: `config.reload()`
|
|
# muta STATE_DIR (testele o folosesc ca sa scoata totul din ~/.claude-discord).
|
|
def bot_log() -> Path:
|
|
return config.LOG_DIR / "bot.log"
|
|
|
|
|
|
def infra_log() -> Path:
|
|
return config.LOG_DIR / "infra.log"
|
|
|
|
|
|
COOKIE_NAME = "dashboard"
|
|
COOKIE_MAX_AGE = 60 * 60 * 24 * 30
|
|
|
|
|
|
def mount_prefix() -> str:
|
|
"""Prefixul sub care e montat panoul (`DASHBOARD_PREFIX`, ex. `/claude`).
|
|
|
|
`tailscale serve --set-path /claude` TAIE prefixul inainte de a proxa, deci in
|
|
mod normal aici ajunge `/api/status`. Prefixul e acceptat totusi si intact,
|
|
pentru cazul unui proxy care nu taie si pentru `curl` direct pe localhost.
|
|
Paginile nu depind de el: toate URL-urile din HTML sunt relative.
|
|
"""
|
|
pfx = (config.get("DASHBOARD_PREFIX") or "").strip().rstrip("/")
|
|
if pfx and not pfx.startswith("/"):
|
|
pfx = "/" + pfx
|
|
return pfx
|
|
|
|
|
|
def strip_prefix(path: str) -> str:
|
|
pfx = mount_prefix()
|
|
if pfx and (path == pfx or path.startswith(pfx + "/")):
|
|
return path[len(pfx):] or "/"
|
|
return path
|
|
|
|
_TOKEN: str | None = None
|
|
|
|
|
|
def reset_token_cache() -> None:
|
|
"""Uita tokenul memorat (folosit de teste dupa `config.reload`)."""
|
|
global _TOKEN
|
|
_TOKEN = None
|
|
|
|
|
|
def auth_disabled() -> bool:
|
|
"""`DASHBOARD_AUTH=off` in `~/.claude-discord/env` scoate complet login-ul.
|
|
|
|
Alegere constienta a operatorului, nu implicit: panoul poate opri un agent cu
|
|
`bypassPermissions` si chei SSH catre tot clusterul, deci fara token oricine
|
|
ajunge la port il poate folosi. Are sens doar pentru ca serviciul e legat de
|
|
`127.0.0.1` si publicat exclusiv in tailnet (`tailscale serve`, tainet only),
|
|
unde identitatea o face deja Tailscale.
|
|
"""
|
|
return (config.get("DASHBOARD_AUTH") or "").strip().lower() in ("off", "none", "0", "false")
|
|
|
|
|
|
def dashboard_token() -> str:
|
|
"""Tokenul de acces, din `~/.claude-discord/env` (`DASHBOARD_TOKEN`).
|
|
|
|
Lipsa lui NU deschide dashboard-ul: se genereaza unul aleator per proces si se
|
|
tipareste in log, deci ramane accesibil doar cui poate citi logul. Ca sa fie
|
|
deschis intentionat, se pune `DASHBOARD_AUTH=off`.
|
|
"""
|
|
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
|
|
|
|
|
|
def tailnet_user(headers) -> str:
|
|
"""Cine e, dupa Tailscale. `tailscale serve` pune antetul pe cererile din
|
|
tailnet; pe localhost lipseste. Doar pentru jurnal — nu e folosit ca decizie
|
|
de acces, fiindca un proces local ar putea sa-l fabrice."""
|
|
return (headers.get("Tailscale-User-Login") or "").strip()
|
|
|
|
|
|
# ── 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:
|
|
"""Secunde de la ultima pornire a unitatii.
|
|
|
|
Sursa principala e `ActiveEnterTimestampMonotonic` (microsecunde pe
|
|
CLOCK_MONOTONIC), comparat cu `time.monotonic()` — ACELASI ceas. **Nu** se
|
|
foloseste `/proc/uptime`: intr-un LXC acela e virtualizat de lxcfs si arata
|
|
uptime-ul containerului, mai mic decat monotonic-ul gazdei pe care il
|
|
raporteaza systemd, deci diferenta iese negativa si uptime-ul apare 0.
|
|
|
|
Rezerva e `ActiveEnterTimestamp`, ora de perete in fusul local al masinii cu
|
|
numele fusului la coada; il taiem si lasam `.timestamp()` sa-l interpreteze
|
|
ca ora locala (`%Z` in strptime nu produce un offset utilizabil).
|
|
"""
|
|
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()
|
|
# "Sun 2026-08-30 12:47:20 UTC" -> data + ora, fara ziua si fusul
|
|
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,
|
|
}
|
|
|
|
|
|
# ── starea punti ────────────────────────────────────────────────────────
|
|
def read_state() -> dict:
|
|
"""state.json fara lock: dashboard-ul doar citeste si nu are voie sa blocheze
|
|
botul. Un JSON prins la mijlocul unei scrieri intoarce {} — se reincarca la
|
|
urmatorul poll."""
|
|
try:
|
|
data = json.loads(config.STATE_FILE.read_text(encoding="utf-8"))
|
|
return data if isinstance(data, dict) else {}
|
|
except (OSError, ValueError):
|
|
return {}
|
|
|
|
|
|
def _pid_alive(pid) -> bool:
|
|
try:
|
|
return pid is not None and Path(f"/proc/{int(pid)}").exists()
|
|
except (TypeError, ValueError):
|
|
return False
|
|
|
|
|
|
def threads_view(state: dict) -> list[dict]:
|
|
out = []
|
|
for tid, rec in (state.get("threads") or {}).items():
|
|
if not isinstance(rec, dict):
|
|
continue
|
|
inflight = rec.get("inflight") or None
|
|
out.append({
|
|
"thread_id": tid,
|
|
"cwd": rec.get("cwd"),
|
|
"model": rec.get("model"),
|
|
"pid": rec.get("pid"),
|
|
"alive": _pid_alive(rec.get("pid")),
|
|
"inflight": bool(inflight),
|
|
"inflight_since": (inflight or {}).get("started_at"),
|
|
"cost_usd": round(float(rec.get("cost_usd_total") or 0), 4),
|
|
"last_active": rec.get("last_active"),
|
|
})
|
|
out.sort(key=lambda t: t.get("last_active") or 0, reverse=True)
|
|
return out
|
|
|
|
|
|
def inflight_threads(state: dict) -> list[str]:
|
|
return [t["thread_id"] for t in threads_view(state) if t["inflight"]]
|
|
|
|
|
|
def pending_approvals() -> list[dict]:
|
|
"""Cererile de confirmare in asteptare, citite direct din director.
|
|
|
|
Nu importam `security.approvals` (API-ul lui e async si porneste un watcher);
|
|
formatul fisierului e fixat in security/README.md.
|
|
"""
|
|
out = []
|
|
d = config.APPROVALS_DIR
|
|
try:
|
|
files = sorted(d.glob("*.json"))
|
|
except OSError:
|
|
return out
|
|
now = time.time()
|
|
for f in files:
|
|
try:
|
|
req = json.loads(f.read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
continue
|
|
if not isinstance(req, dict) or req.get("status") != "pending":
|
|
continue
|
|
out.append({
|
|
"request_id": req.get("request_id"),
|
|
"thread_id": req.get("thread_id"),
|
|
"tool_name": req.get("tool_name"),
|
|
"command": (req.get("command") or "")[:500],
|
|
"rule": req.get("rule"),
|
|
"reason": req.get("reason"),
|
|
"created_at": req.get("created_at"),
|
|
"expires_in": (
|
|
round(req["expires_at"] - now, 1)
|
|
if isinstance(req.get("expires_at"), (int, float)) else None
|
|
),
|
|
})
|
|
return out
|
|
|
|
|
|
def doctor() -> list[dict]:
|
|
checks: list[dict] = []
|
|
|
|
info = unit_info(SERVICE)
|
|
checks.append({
|
|
"name": "Serviciu claude-discord",
|
|
"pass": info["active"],
|
|
"detail": f'{info["state"]}/{info["sub"]}, {info["restarts"]} restarturi',
|
|
})
|
|
|
|
st = read_state()
|
|
checks.append({
|
|
"name": "state.json",
|
|
"pass": bool(st),
|
|
"detail": f'{len(st.get("threads") or {})} fire' if st else "ilizibil sau gol",
|
|
})
|
|
|
|
cap = parse_cap(config.get("COST_CAP_USD_DAY"), 0.0)
|
|
spent = float((st.get("cost") or {}).get("usd") or 0)
|
|
checks.append({
|
|
"name": "Plafon de cost pe zi",
|
|
"pass": cap <= 0 or spent < cap,
|
|
"detail": f"{spent:.2f} / {cap:.2f} USD" if cap > 0
|
|
else f"{spent:.2f} USD cheltuiti, fara plafon (abonament)",
|
|
})
|
|
|
|
try:
|
|
du = shutil.disk_usage("/")
|
|
pct = du.free / du.total * 100
|
|
checks.append({
|
|
"name": "Spatiu pe disc",
|
|
"pass": pct > 10,
|
|
"detail": f"{pct:.1f}% liber ({du.free // 1024**3} GB)",
|
|
})
|
|
except OSError as exc:
|
|
checks.append({"name": "Spatiu pe disc", "pass": False, "detail": str(exc)})
|
|
|
|
try:
|
|
log = bot_log()
|
|
size_mb = log.stat().st_size / 1024**2 if log.exists() else 0
|
|
checks.append({
|
|
"name": "bot.log",
|
|
"pass": log.exists() and size_mb < 100,
|
|
"detail": f"{size_mb:.1f} MB" if log.exists() else "lipseste",
|
|
})
|
|
except OSError as exc:
|
|
checks.append({"name": "bot.log", "pass": False, "detail": str(exc)})
|
|
|
|
# Regula deny(ssh) a mai taiat o data accesul la infrastructura — vezi
|
|
# security/README.md. Verificam sa nu reapara la o editare viitoare.
|
|
try:
|
|
settings = json.loads(config.SETTINGS_FILE.read_text(encoding="utf-8"))
|
|
deny = (settings.get("permissions") or {}).get("deny") or []
|
|
bad = [d for d in deny if d.startswith(("Bash(ssh", "Bash(scp"))]
|
|
checks.append({
|
|
"name": "Reguli deny in bot-settings.json",
|
|
"pass": not bad,
|
|
"detail": f"blocheaza infrastructura: {bad}" if bad else f"{len(deny)} reguli, ssh liber",
|
|
})
|
|
except (OSError, ValueError) as exc:
|
|
checks.append({"name": "Reguli deny in bot-settings.json", "pass": False, "detail": str(exc)})
|
|
|
|
hook = _BRIDGE / "security" / "confirm_hook.py"
|
|
checks.append({
|
|
"name": "Hook de confirmare",
|
|
"pass": hook.exists(),
|
|
"detail": str(hook) if hook.exists() else "lipseste",
|
|
})
|
|
|
|
claude_bin = shutil.which(config.get("CLAUDE_BIN") or "claude")
|
|
checks.append({
|
|
"name": "CLI claude",
|
|
"pass": bool(claude_bin),
|
|
"detail": claude_bin or "nu e in PATH",
|
|
})
|
|
|
|
return checks
|
|
|
|
|
|
def orphans_report(dry_run: bool = True) -> dict:
|
|
"""`/cleanup` din Discord, expus si aici. Importul e lenes fiindca modulul
|
|
citeste /proc la import-time in unele cai."""
|
|
import cleanup # noqa: PLC0415
|
|
|
|
found = cleanup.find_orphans(read_state())
|
|
results = None
|
|
if not dry_run and found:
|
|
results = cleanup.kill_orphans(found, dry_run=False)
|
|
return {
|
|
"orphans": found,
|
|
"killed": results,
|
|
"report": cleanup.format_report(found, results),
|
|
}
|
|
|
|
|
|
# ── 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 = "claude-discord-dashboard"
|
|
protocol_version = "HTTP/1.1"
|
|
|
|
def __init__(self, *a, **kw):
|
|
super().__init__(*a, directory=str(_DASH), **kw)
|
|
|
|
# --- utilitare -----------------------------------------------------
|
|
def log_message(self, format, *args): # jurnal compact, o linie # 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):
|
|
raw = urlparse(self.path).path
|
|
path = strip_prefix(raw)
|
|
if path == "/" and not raw.endswith("/"):
|
|
# `/claude` fara slash final: URL-urile relative din pagina s-ar
|
|
# rezolva la radacina hostului (`/api/status`), unde proxy-ul nu mai
|
|
# trimite nimic incoace. Fortam forma cu slash.
|
|
self.send_response(301)
|
|
self.send_header("Location", raw + "/")
|
|
self.send_header("Content-Length", "0")
|
|
self.end_headers()
|
|
return
|
|
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():
|
|
# Prefixul reintra AICI in mod deliberat. Un "/login.html" absolut ar
|
|
# arunca browserul in radacina hostului (alt serviciu), iar un
|
|
# "login.html" relativ se rezolva gresit cand adresa vine fara slash
|
|
# final (`/claude` -> `/login.html`). Cu prefixul reatasat, ambele
|
|
# forme ajung unde trebuie, si direct pe localhost la fel: calea de
|
|
# intrare e curatata oricum de `strip_prefix`.
|
|
self.send_response(302)
|
|
self.send_header("Location", mount_prefix() + "/login.html")
|
|
self.send_header("Content-Length", "0")
|
|
self.end_headers()
|
|
return
|
|
if path in ("/", "/index.html"):
|
|
return self.send_html("index.html")
|
|
if path == "/login.html":
|
|
if auth_disabled():
|
|
self.send_response(302)
|
|
self.send_header("Location", mount_prefix() + "/")
|
|
self.send_header("Content-Length", "0")
|
|
self.end_headers()
|
|
return None
|
|
return self.send_html("login.html")
|
|
self.path = path
|
|
return super().do_GET()
|
|
|
|
def send_html(self, name: str):
|
|
"""Trimite o pagina, cu `<base href>` pus la servire.
|
|
|
|
De ce e nevoie: `tailscale serve --set-path /claude` TAIE prefixul, deci
|
|
serverul nu poate sti daca browserul e la `/claude` sau la `/claude/` —
|
|
ambele ajung aici ca `/`. Fara slash final, un URL relativ (`api/status`)
|
|
se rezolva la radacina hostului, unde proxy-ul nu mai trimite nimic
|
|
incoace: pagina se incarca si ramane goala, fara nicio eroare vizibila.
|
|
`<base href="/claude/">` fixeaza rezolvarea indiferent de forma adresei.
|
|
"""
|
|
try:
|
|
body = (_DASH / name).read_text(encoding="utf-8")
|
|
except OSError:
|
|
return self.send_error(404)
|
|
pfx = mount_prefix()
|
|
if pfx:
|
|
body = body.replace("<head>", f'<head>\n<base href="{pfx}/">', 1)
|
|
data = body.encode("utf-8")
|
|
self.send_response(200)
|
|
self.send_header("Content-Type", "text/html; charset=utf-8")
|
|
self.send_header("Content-Length", str(len(data)))
|
|
# no-store: paginile poarta acum si configuratia (base href), deci o
|
|
# copie veche din cache ar trimite cererile in alta parte.
|
|
self.send_header("Cache-Control", "no-store")
|
|
self.end_headers()
|
|
try:
|
|
self.wfile.write(data)
|
|
except (BrokenPipeError, ConnectionResetError):
|
|
pass
|
|
|
|
def route_get(self, path: str):
|
|
qs = parse_qs(urlparse(self.path).query)
|
|
if path == "/api/status":
|
|
state = read_state()
|
|
return self.send_json({
|
|
"service": unit_info(SERVICE),
|
|
"dashboard": unit_info(SELF_SERVICE),
|
|
"threads": threads_view(state),
|
|
"cost": state.get("cost") or {},
|
|
"cost_cap": parse_cap(config.get("COST_CAP_USD_DAY"), 0.0),
|
|
"pending_approvals": len(pending_approvals()),
|
|
"auth": not auth_disabled(),
|
|
"user": tailnet_user(self.headers),
|
|
"now": time.time(),
|
|
})
|
|
if path == "/api/logs":
|
|
try:
|
|
n = min(max(int(qs.get("lines", ["200"])[0]), 1), 2000)
|
|
except ValueError:
|
|
n = 200
|
|
which = qs.get("file", ["bot"])[0]
|
|
target = infra_log() if which == "infra" else bot_log()
|
|
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()})
|
|
if path == "/api/approvals":
|
|
return self.send_json({"approvals": pending_approvals()})
|
|
if path == "/api/doctor":
|
|
return self.send_json({"checks": doctor()})
|
|
if path == "/api/cleanup":
|
|
return self.send_json(orphans_report(dry_run=True))
|
|
if path == "/api/maria/status":
|
|
return self.send_json({
|
|
"bridge": unit_info(MARIA_UNITS["bridge"]),
|
|
"rag": unit_info(MARIA_UNITS["rag"]),
|
|
"whatsapp": maria_whatsapp_status(),
|
|
"index": maria_index_info(),
|
|
"documents": len(maria_documents()),
|
|
"sync": {**maria_sync_state(), "drive_remote": maria_get("DRIVE_REMOTE") or None},
|
|
"support_jid": maria_get("SUPPORT_JID") or None,
|
|
})
|
|
if path == "/api/maria/documents":
|
|
return self.send_json({"documents": maria_documents()})
|
|
if path == "/api/maria/escalations":
|
|
return self.send_json({"escalations": maria_escalations()})
|
|
if path == "/api/maria/clients":
|
|
return self.send_json({"clients": maria_clients()})
|
|
if path == "/api/maria/logs":
|
|
try:
|
|
n = min(max(int(qs.get("lines", ["200"])[0]), 1), 2000)
|
|
except ValueError:
|
|
n = 200
|
|
target = maria_log(qs.get("service", ["bridge"])[0])
|
|
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 = strip_prefix(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/cleanup":
|
|
data = self.read_json()
|
|
return self.send_json(orphans_report(dry_run=bool(data.get("dry_run", True))))
|
|
if path == "/api/approvals/decide":
|
|
return self.handle_decide()
|
|
if path == "/api/maria/service":
|
|
return self.handle_maria_service()
|
|
if path == "/api/maria/pair":
|
|
return self.handle_maria_pair()
|
|
if path == "/api/maria/documents":
|
|
return self.handle_maria_document_write()
|
|
if path == "/api/maria/documents/delete":
|
|
return self.handle_maria_document_delete()
|
|
if path == "/api/maria/clients":
|
|
return self.handle_maria_client_write()
|
|
if path == "/api/maria/clients/delete":
|
|
return self.handle_maria_client_delete()
|
|
if path == "/api/maria/reindex":
|
|
return self.send_json(maria_reindex())
|
|
if path == "/api/maria/sync":
|
|
data = self.read_json()
|
|
return self.send_json(maria_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), 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) # incetineste ghicitul
|
|
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 UNITATEA FIXA. Numele nu vine din request."""
|
|
data = self.read_json()
|
|
action = str(data.get("action") or "")
|
|
if action not in ("start", "stop", "restart"):
|
|
return self.send_json({"ok": False, "error": f"actiune necunoscuta: {action}"}, 400)
|
|
|
|
if action in ("stop", "restart") and not data.get("force"):
|
|
busy = inflight_threads(read_state())
|
|
if busy:
|
|
return self.send_json({
|
|
"ok": False,
|
|
"error": "tururi in desfasurare",
|
|
"inflight": busy,
|
|
"hint": "retrimite cu force=true ca sa le intrerupi",
|
|
}, 409)
|
|
|
|
who = tailnet_user(self.headers) or "local"
|
|
log_line = f"[actiune] {action} pe {SERVICE}, cerut de {who}"
|
|
print(log_line, file=sys.stderr, flush=True)
|
|
try:
|
|
r = _sysctl(action, SERVICE)
|
|
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) # lasa systemd sa actualizeze starea inainte de raspuns
|
|
return self.send_json({"ok": True, "action": action, "service": unit_info(SERVICE)})
|
|
|
|
def handle_decide(self):
|
|
data = self.read_json()
|
|
rid = str(data.get("request_id") or "")
|
|
decision = str(data.get("decision") or "")
|
|
# `allow_session` = permite si nu mai intreba in firul asta (vezi
|
|
# security/README.md, sectiunea "Aprobari valabile pe tot firul").
|
|
if decision not in ("allow", "allow_session", "deny"):
|
|
return self.send_json({"ok": False, "error": "decizie invalida"}, 400)
|
|
if not rid or "/" in rid or ".." in rid:
|
|
return self.send_json({"ok": False, "error": "request_id invalid"}, 400)
|
|
import importlib # noqa: PLC0415
|
|
approvals = importlib.import_module("security.approvals")
|
|
ok = approvals.submit_decision(rid, decision)
|
|
return self.send_json({"ok": bool(ok), "request_id": rid, "decision": decision},
|
|
200 if ok else 404)
|
|
|
|
def handle_maria_service(self):
|
|
"""start / stop / restart pe una din UNITATILE FIXE ale Mariei — acelasi
|
|
principiu ca handle_service: numele vine dintr-o cheie (bridge/rag), nu
|
|
dintr-un nume de unit arbitrar din request."""
|
|
data = self.read_json()
|
|
which = str(data.get("service") or "")
|
|
action = str(data.get("action") or "")
|
|
unit = MARIA_UNITS.get(which)
|
|
if not unit:
|
|
return self.send_json({"ok": False, "error": f"serviciu Maria necunoscut: {which}"}, 400)
|
|
if action not in ("start", "stop", "restart"):
|
|
return self.send_json({"ok": False, "error": f"actiune necunoscuta: {action}"}, 400)
|
|
who = tailnet_user(self.headers) or "local"
|
|
print(f"[actiune] {action} pe {unit}, cerut de {who}", 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_maria_pair(self):
|
|
"""Cod de asociere WhatsApp — alternativa la QR, cand nu poti scana ecranul."""
|
|
data = self.read_json()
|
|
phone = str(data.get("phone") or "")
|
|
who = tailnet_user(self.headers) or "local"
|
|
print(f"[actiune] cod de asociere WhatsApp pentru {phone}, cerut de {who}",
|
|
file=sys.stderr, flush=True)
|
|
out, status = maria_request_pairing_code(phone)
|
|
return self.send_json(out, status)
|
|
|
|
def handle_maria_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:
|
|
maria_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_maria_document_delete(self):
|
|
data = self.read_json()
|
|
name = str(data.get("name") or "")
|
|
try:
|
|
ok = maria_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 handle_maria_client_write(self):
|
|
data = self.read_json()
|
|
cid = str(data.get("id") or "")
|
|
numere = data.get("numere") or []
|
|
if isinstance(numere, str):
|
|
numere = numere.splitlines()
|
|
try:
|
|
maria_save_client(cid, str(data.get("nume") or ""), str(data.get("profil") or ""),
|
|
list(numere))
|
|
except (ValueError, OSError) as exc:
|
|
return self.send_json({"ok": False, "error": str(exc)}, 400)
|
|
return self.send_json({"ok": True, "id": cid})
|
|
|
|
def handle_maria_client_delete(self):
|
|
data = self.read_json()
|
|
cid = str(data.get("id") or "")
|
|
ok = maria_delete_client(cid)
|
|
return self.send_json({"ok": ok, "id": cid}, 200 if ok else 404)
|
|
|
|
|
|
def main() -> None:
|
|
bind = config.get("DASHBOARD_BIND") or "127.0.0.1"
|
|
try:
|
|
port = int(config.get("DASHBOARD_PORT") or 18790)
|
|
except ValueError:
|
|
port = 18790
|
|
dashboard_token() # forteaza avertismentul de token la pornire, nu la primul GET
|
|
srv = ThreadingHTTPServer((bind, port), Handler)
|
|
srv.daemon_threads = True
|
|
print(f"dashboard pe http://{bind}:{port} (unitate controlata: {SERVICE})",
|
|
file=sys.stderr, flush=True)
|
|
try:
|
|
srv.serve_forever()
|
|
except KeyboardInterrupt:
|
|
pass
|
|
finally:
|
|
srv.server_close()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|