"""Proxy VM112 domínios orquestrados + limpeza Desk (Spec 017).""" from __future__ import annotations import os import sqlite3 import time from datetime import datetime, timezone from threading import Lock from typing import Any import httpx from app import auth VM112_API = os.getenv("VM112_API_URL", "http://10.10.10.112:8090") VM112_ADMIN_API_KEY = os.getenv("VM112_ADMIN_API_KEY", "ibytera-corp-api-key-change-later") # Purge Carbonio/CF/Traefik pode demorar vários minutos — evitar httpx "timed out" prematuro. VM112_PURGE_HTTP_TIMEOUT = float(os.getenv("VM112_PURGE_HTTP_TIMEOUT", "300")) VM112_LIST_TIMEOUT = float(os.getenv("VM112_LIST_TIMEOUT", "30")) VM112_DOMAINS_CACHE_TTL = float(os.getenv("VM112_DOMAINS_CACHE_TTL", "60")) VM112_DOMAIN_DETAIL_TTL = float(os.getenv("VM112_DOMAIN_DETAIL_TTL", "45")) _DOMAINS_CACHE: dict[str, Any] | None = None _DOMAINS_CACHE_AT = 0.0 _DOMAINS_CACHE_LOCK = Lock() _DOMAIN_DETAIL_CACHE: dict[str, tuple[float, dict[str, Any]]] = {} PURGE_BLOCKLIST = frozenset({"ligbox.com.br", "itecnologys.com"}) # Exige código de autorização gerado em Infra (root) — além de senha Root no purge. PURGE_EXTRA_AUTH_DOMAINS = frozenset({"myvexx.com"}) VM112_PURGE_STEP_LABELS = ( "Contas Carbonio (zmprov da)", "Domínio Carbonio (zmprov dd)", "Portal users Self-Service", "Pasta ligbox-sites", "Zona Cloudflare Ibytera", "Traefik / SNI CT114", "Logs de sessão wizard", ) # Origem da solicitação de expurgo (user = JWT / sessão; source = canal) PURGE_SOURCE_LABELS: dict[str, str] = { "desk.ui.purge": "Desk › Purge domínio (UI)", "desk.api.purge": "Desk › Purge domínio (API sync)", "desk.api.purge-jobs": "Desk › Purge domínio (job async)", "desk.api.purge-stream": "Desk › Purge domínio (SSE)", "cursor.agent": "Cursor Agent (comando chat)", "domain-console.sandbox": "Domain Console › sandbox", "desk.agentic": "Desk › rotina agentic", } def purge_source_label(source: str) -> str: return PURGE_SOURCE_LABELS.get(source, source or "desconhecido") def _ts() -> str: return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") def _timeline_entry(label: str, status: str, detail: str = "") -> dict[str, str]: return {"at": _ts(), "label": label, "status": status, "detail": detail} def _vm112_headers() -> dict[str, str]: return {"X-Api-Key": VM112_ADMIN_API_KEY} def verify_root_password(conn: sqlite3.Connection, password: str) -> bool: row = conn.execute( "SELECT password_hash FROM desk_users WHERE username = 'root' AND active = 1" ).fetchone() if not row or not row["password_hash"]: return False return auth.verify_password(password, row["password_hash"]) def requires_purge_extra_auth(domain: str) -> bool: return domain.lower().strip() in PURGE_EXTRA_AUTH_DOMAINS def delete_carbonio_account(email: str) -> dict[str, Any]: """Remove uma conta Carbonio (zmprov da) — Spec 022.""" email = email.lower().strip() if "@" not in email: raise ValueError("e-mail inválido") domain = email.split("@", 1)[1] if domain in PURGE_BLOCKLIST: raise ValueError(f"Domínio protegido: {domain}") with httpx.Client(timeout=120.0) as client: r = client.post( f"{VM112_API}/api/admin/accounts/{email}/delete", headers=_vm112_headers(), ) if r.status_code == 404: return {"ok": True, "email": email, "message": "Conta já não existia no Carbonio", "skipped": True} r.raise_for_status() data = r.json() return { "ok": True, "email": email, "message": data.get("message") or f"Conta {email} removida", "detail": data, } def list_domains(query: str = "", *, force_refresh: bool = False) -> dict[str, Any]: """Lista domínios VM112 — cache TTL curto (Serviços IaaS poll frequente).""" global _DOMAINS_CACHE, _DOMAINS_CACHE_AT if not query and not force_refresh: with _DOMAINS_CACHE_LOCK: if _DOMAINS_CACHE is not None and time.time() - _DOMAINS_CACHE_AT < VM112_DOMAINS_CACHE_TTL: return {**_DOMAINS_CACHE, "cached": True, "cache_age_sec": int(time.time() - _DOMAINS_CACHE_AT)} with httpx.Client(timeout=VM112_LIST_TIMEOUT) as client: r = client.get( f"{VM112_API}/api/admin/domains", params={"q": query} if query else None, headers=_vm112_headers(), ) r.raise_for_status() data = r.json() if not query: with _DOMAINS_CACHE_LOCK: _DOMAINS_CACHE = data _DOMAINS_CACHE_AT = time.time() return data def get_domain(domain: str) -> dict[str, Any]: domain = domain.lower().strip() now = time.time() cached = _DOMAIN_DETAIL_CACHE.get(domain) if cached and now - cached[0] < VM112_DOMAIN_DETAIL_TTL: data = dict(cached[1]) else: with httpx.Client(timeout=VM112_LIST_TIMEOUT) as client: r = client.get( f"{VM112_API}/api/admin/domains/{domain}", headers=_vm112_headers(), ) r.raise_for_status() data = r.json() if isinstance(data, dict): _DOMAIN_DETAIL_CACHE[domain] = (now, data) if isinstance(data, dict): data["purge_extra_auth_required"] = requires_purge_extra_auth(domain) data["purge_blocked"] = domain in PURGE_BLOCKLIST return data def domain_exists_on_vm112(domain: str) -> bool: """True se o domínio ainda consta na lista orquestrada VM112.""" domain = domain.lower().strip() try: data = list_domains() items = data.get("domains") if isinstance(data, dict) else data if not isinstance(items, list): return True for item in items: name = item.get("domain") if isinstance(item, dict) else item if str(name or "").lower().strip() == domain: return True return False except Exception: # VM112 indisponível — não assumir removido durante poll return True def start_purge_vm112(domain: str) -> dict[str, Any]: """Inicia purge assíncrono na VM112 (Spec 017 Fase 3).""" domain = domain.lower().strip() with httpx.Client(timeout=VM112_PURGE_HTTP_TIMEOUT) as client: r = client.post( f"{VM112_API}/api/admin/domains/{domain}/purge", headers=_vm112_headers(), ) r.raise_for_status() return r.json() def poll_purge_vm112_job(job_id: str) -> dict[str, Any]: with httpx.Client(timeout=VM112_PURGE_HTTP_TIMEOUT) as client: r = client.get( f"{VM112_API}/api/admin/domains/purge-jobs/{job_id}", headers=_vm112_headers(), ) r.raise_for_status() return r.json() def vm112_job_steps_timeline(job: dict[str, Any]) -> list[dict[str, str]]: """Passos individuais VM112 durante execução (Fase 3).""" out: list[dict[str, str]] = [] for step in job.get("steps") or []: if not isinstance(step, dict): continue st = str(step.get("status") or "pending") if st == "pending": continue label = str(step.get("label") or "Passo VM112") if st == "done": status = "ok" elif st == "error": status = "fail" else: status = "running" detail = str(step.get("detail") or "") at = step.get("finished_at") or step.get("started_at") or _ts() out.append({"at": at, "label": label, "status": status, "detail": detail}) return out def purge_vm112_with_poll(domain: str, poll_interval: float = 1.5, timeout: float = 600.0): """Generator: (event_type, payload) — passos em tempo real + resultado final.""" import time started = start_purge_vm112(domain) job_id = started.get("job_id") if not job_id: yield ("final", started) return t0 = time.monotonic() deadline = t0 + timeout seen = 0 poll_errors = 0 while time.monotonic() < deadline: try: job = poll_purge_vm112_job(job_id) poll_errors = 0 except Exception as exc: poll_errors += 1 err = str(exc) or "erro poll VM112" yield ( "heartbeat", { "elapsed": int(time.monotonic() - t0), "job_id": job_id, "warn": err, "poll_errors": poll_errors, }, ) if poll_errors >= 5 and not domain_exists_on_vm112(domain): yield ( "final", { "ok": True, "job_id": job_id, "recovered": True, "steps": [], "result": {"message": "Domínio ausente na VM112 após falhas de poll"}, }, ) return time.sleep(poll_interval) continue steps = vm112_job_steps_timeline(job) if len(steps) > seen: for step in steps[seen:]: yield ("step", step) seen = len(steps) status = job.get("status") if status == "completed": yield ( "final", { "ok": True, "job_id": job_id, "steps": steps, "result": job.get("result") or {}, }, ) return if status == "failed": yield ( "final", { "ok": False, "job_id": job_id, "steps": steps, "error": job.get("error") or "Purge VM112 falhou", "result": job.get("result") or {}, }, ) return yield ("heartbeat", {"elapsed": int(time.monotonic() - t0), "job_id": job_id}) time.sleep(poll_interval) if not domain_exists_on_vm112(domain): yield ( "final", { "ok": True, "job_id": job_id, "recovered": True, "steps": [], "result": {"message": "Domínio removido na VM112 (timeout poll — recuperado)"}, }, ) return yield ( "final", { "ok": False, "error": "Timeout purge VM112 — domínio ainda presente; use Recuperar ou repita", "job_id": job_id, }, ) def purge_vm112(domain: str) -> dict[str, Any]: domain = domain.lower().strip() for kind, payload in purge_vm112_with_poll(domain): if kind == "final": return payload return {"ok": False, "error": "Purge VM112 sem resposta"} def vm112_purge_timeline(vm112_result: dict[str, Any]) -> list[dict[str, str]]: """Converte resposta VM112 em linhas de timeline.""" raw_steps = vm112_result.get("steps") if isinstance(raw_steps, list) and raw_steps: out: list[dict[str, str]] = [] for step in raw_steps: if not isinstance(step, dict): continue label = str(step.get("label") or step.get("name") or "Passo VM112") ok = step.get("ok", step.get("success", True)) status = "ok" if ok else "fail" detail = str(step.get("message") or step.get("detail") or "") at = step.get("at") or _ts() out.append({"at": at, "label": label, "status": status, "detail": detail}) return out if vm112_result.get("ok") is False: return [ _timeline_entry( "Purge VM112", "fail", str(vm112_result.get("message") or vm112_result.get("error") or "falhou"), ) ] return [_timeline_entry("Purge VM112", "ok", "Orquestração VM112 concluída")] def purge_desk_records(conn: sqlite3.Connection, domain: str) -> dict[str, int]: counts, _ = purge_desk_timeline( conn, domain, by_user="system", purge_source="desk.api.purge", record_event=False ) return counts def purge_domain_console_scenario(conn: sqlite3.Connection, domain: str) -> int: domain = domain.lower().strip() return conn.execute( "DELETE FROM domain_console_scenarios WHERE domain = ?", (domain,) ).rowcount def record_domain_purged_event( conn: sqlite3.Connection, domain: str, *, by_user: str, purge_source: str, desk_removed: dict[str, int] | None = None, vm112_ok: bool | None = None, vm112_error: str | None = None, job_id: str | None = None, ) -> int: """Marca expurgo em eventos — histórico preservado (não apaga webhook_events).""" import json now = datetime.now(timezone.utc).isoformat() payload = json.dumps( { "domain": domain.lower().strip(), "event": "domain.purged", "purged_at": now, "by_user": by_user, "purge_source": purge_source, "purge_tool": purge_source_label(purge_source), "message": ( f"Domínio expurgado em {now} pelo utilizador {by_user} " f"via {purge_source_label(purge_source)}" ), "desk_removed": desk_removed or {}, "vm112_ok": vm112_ok, "vm112_error": vm112_error, "job_id": job_id, "accounts_hint": [f"admin@{domain.lower().strip()}", f"mail.{domain.lower().strip()}"], }, ensure_ascii=False, ) cur = conn.execute( "INSERT INTO webhook_events (event_type, source, payload, created_at) VALUES (?,?,?,?)", ("domain.purged", "desk.purge", payload, now), ) conn.commit() return int(cur.lastrowid) def purge_desk_timeline( conn: sqlite3.Connection, domain: str, *, by_user: str = "system", purge_source: str = "desk.api.purge", record_event: bool = True, vm112_ok: bool | None = None, vm112_error: str | None = None, job_id: str | None = None, ) -> tuple[dict[str, int], list[dict[str, str]]]: """Purge Desk — mantém webhook_events; regista domain.purged no fim.""" domain = domain.lower().strip() like = f"%{domain}%" timeline: list[dict[str, str]] = [] counts: dict[str, int] = {} desk_steps = ( ("Desk — domain_console_scenarios", "domain_console_scenarios", None, None), ("Desk — tickets", "tickets", "DELETE FROM tickets WHERE subject LIKE ? OR payload LIKE ?", (like, like)), ("Desk — audit_domains", "audit_domains", "DELETE FROM audit_domains WHERE domain = ?", (domain,)), ("Desk — assist_sessions", "assist_sessions", "DELETE FROM assist_sessions WHERE domain = ?", (domain,)), ("Desk — audit_checks", "audit_checks", "DELETE FROM audit_checks WHERE domain = ?", (domain,)), ) for label, key, sql, params in desk_steps: if key == "domain_console_scenarios": n = purge_domain_console_scenario(conn, domain) else: n = conn.execute(sql, params).rowcount counts[key] = n timeline.append(_timeline_entry(label, "ok", f"{n} registo(s) removido(s)")) event_id = 0 if record_event: event_id = record_domain_purged_event( conn, domain, by_user=by_user, purge_source=purge_source, desk_removed=counts, vm112_ok=vm112_ok, vm112_error=vm112_error, job_id=job_id, ) counts["webhook_events_purged_marker"] = 1 timeline.append( _timeline_entry( "Desk — evento domain.purged", "ok", f"event_id={event_id} · histórico webhook preservado", ) ) conn.commit() return counts, timeline def build_purge_timeline(vm112_result: dict[str, Any], desk_counts: dict[str, int], desk_timeline: list[dict[str, str]]) -> list[dict[str, str]]: timeline = [_timeline_entry("Validação Root + confirmação", "ok")] timeline.extend(vm112_purge_timeline(vm112_result)) timeline.extend(desk_timeline) total_desk = sum(desk_counts.values()) timeline.append(_timeline_entry("Purge concluído", "ok", f"Desk: {total_desk} registo(s)")) return timeline