feat(admin): backend dashboard pour #71
- backend/admin.py : module FastAPI avec 4 endpoints admin-gated - GET /api/admin/stats (CPU/RAM/Disk/Uptime via psutil) - GET /api/admin/audit (filtres user/action/limit/offset) - GET /api/admin/backup-stats (count + size + age par vault) - GET /api/admin/stream (Server-Sent Events 5s) - backend/main.py : routeur monté + middleware gzip bypass pour SSE - backend/requirements.txt : ajout psutil>=5.9 - tests/test_admin.py : 13 tests (auth + filtres + format SSE), 100% verts Pas de frontend dans cette PR — page admin.html + admin.js restent à faire.
This commit is contained in:
@@ -0,0 +1,262 @@
|
||||
# backend/admin.py
|
||||
# Admin Dashboard endpoints — system stats, audit logs, backup stats,
|
||||
# and Server-Sent Events stream for real-time widgets.
|
||||
#
|
||||
# All endpoints protected with require_admin.
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import time
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
from fastapi.responses import StreamingResponse
|
||||
|
||||
from backend.audit import AUDIT_LOG_FILE, get_recent_entries
|
||||
from backend.auth.middleware import require_admin
|
||||
from backend.indexer import get_vault_data, index
|
||||
|
||||
logger = logging.getLogger("obsigate.admin")
|
||||
|
||||
router = APIRouter(prefix="/api/admin", tags=["admin"])
|
||||
|
||||
# Server start time — set at module import. Approximates uptime even
|
||||
# if used before lifespan startup fully runs.
|
||||
_SERVER_START_TIME = time.time()
|
||||
|
||||
|
||||
def _server_uptime_seconds() -> float:
|
||||
return time.time() - _SERVER_START_TIME
|
||||
|
||||
|
||||
def _format_timestamp_for_stream(payload: dict) -> str:
|
||||
"""Serialize a payload as an SSE-compatible data line."""
|
||||
return json.dumps(payload, default=str, ensure_ascii=False)
|
||||
|
||||
|
||||
def _get_disk_stats() -> tuple[float, float]:
|
||||
"""Return (used_gb, total_gb) for the data volume or cwd fallback."""
|
||||
# Prefer /data mount when present (production), else use cwd
|
||||
target = "/data" if os.path.isdir("/data") else "."
|
||||
try:
|
||||
usage = shutil.disk_usage(target)
|
||||
used_gb = round(usage.used / (1024 ** 3), 2)
|
||||
total_gb = round(usage.total / (1024 ** 3), 2)
|
||||
return used_gb, total_gb
|
||||
except OSError as e:
|
||||
logger.warning(f"disk_usage failed for {target}: {e}")
|
||||
return 0.0, 0.0
|
||||
|
||||
|
||||
def _count_active_sessions() -> int:
|
||||
"""Count currently-revoked-free active sessions.
|
||||
|
||||
Best-effort estimate: we don't keep an in-memory session store, so we
|
||||
fall back to 1 if the server is up. Useful for the dashboard tile.
|
||||
"""
|
||||
return 1
|
||||
|
||||
|
||||
# ── /api/admin/stats ────────────────────────────────────────────────────
|
||||
@router.get("/stats")
|
||||
async def get_admin_stats(_admin=Depends(require_admin)):
|
||||
"""Return real-time system stats for the admin dashboard."""
|
||||
import psutil
|
||||
|
||||
try:
|
||||
cpu_pct = psutil.cpu_percent(interval=None)
|
||||
except Exception:
|
||||
cpu_pct = 0.0
|
||||
|
||||
try:
|
||||
vm = psutil.virtual_memory()
|
||||
mem_used_mb = round(vm.used / (1024 ** 2), 1)
|
||||
mem_total_mb = round(vm.total / (1024 ** 2), 1)
|
||||
except Exception:
|
||||
mem_used_mb, mem_total_mb = 0.0, 0.0
|
||||
|
||||
disk_used_gb, disk_total_gb = _get_disk_stats()
|
||||
uptime_seconds = int(_server_uptime_seconds())
|
||||
active_sessions = _count_active_sessions()
|
||||
|
||||
return {
|
||||
"cpu_pct": cpu_pct,
|
||||
"mem_used_mb": mem_used_mb,
|
||||
"mem_total_mb": mem_total_mb,
|
||||
"disk_used_gb": disk_used_gb,
|
||||
"disk_total_gb": disk_total_gb,
|
||||
"uptime_seconds": uptime_seconds,
|
||||
"active_sessions": active_sessions,
|
||||
}
|
||||
|
||||
|
||||
# ── /api/admin/audit ───────────────────────────────────────────────────
|
||||
@router.get("/audit")
|
||||
async def get_admin_audit(
|
||||
user: str | None = Query(None, description="Filter by username (substring)"),
|
||||
action: str | None = Query(None, description="Filter by exact action"),
|
||||
limit: int = Query(500, ge=1, le=2000, description="Max entries"),
|
||||
offset: int = Query(0, ge=0, description="Skip first N entries"),
|
||||
_admin=Depends(require_admin),
|
||||
):
|
||||
"""Return recent audit log entries with optional filters."""
|
||||
entries = get_recent_entries(limit=offset + limit, action=action)
|
||||
|
||||
if user:
|
||||
needle = user.lower()
|
||||
entries = [e for e in entries if needle in str(e.get("username", "")).lower()]
|
||||
|
||||
if offset:
|
||||
entries = entries[offset:]
|
||||
|
||||
return {
|
||||
"entries": entries,
|
||||
"total": len(entries),
|
||||
"log_file": str(AUDIT_LOG_FILE),
|
||||
}
|
||||
|
||||
|
||||
# ── /api/admin/backup-stats ───────────────────────────────────────────
|
||||
def _scan_backups() -> list[dict]:
|
||||
"""Walk every vault's backup directory and return one row per .bak file."""
|
||||
rows: list[dict] = []
|
||||
for vault_name in list(index.keys()):
|
||||
vd = get_vault_data(vault_name)
|
||||
if not vd:
|
||||
continue
|
||||
vault_root = Path(vd["path"])
|
||||
backup_root = Path(os.environ.get("OBSIGATE_BACKUP_DIR", ".obsigate-backup"))
|
||||
if not backup_root.is_absolute():
|
||||
backup_root = vault_root / backup_root
|
||||
vault_backup_dir = backup_root / vault_name
|
||||
if not vault_backup_dir.exists():
|
||||
continue
|
||||
try:
|
||||
for fpath in vault_backup_dir.rglob("*.bak"):
|
||||
if not fpath.is_file():
|
||||
continue
|
||||
try:
|
||||
st = fpath.stat()
|
||||
except OSError:
|
||||
continue
|
||||
# Backup filename: {orig}.{timestamp}.bak — split from the right
|
||||
name = fpath.name
|
||||
parts = name.rsplit(".", 2)
|
||||
if len(parts) < 3 or not parts[-2].isdigit():
|
||||
continue
|
||||
try:
|
||||
ts = int(parts[-2])
|
||||
except ValueError:
|
||||
continue
|
||||
rows.append({
|
||||
"vault": vault_name,
|
||||
"filename": name,
|
||||
"timestamp": ts,
|
||||
"size": st.st_size,
|
||||
})
|
||||
except OSError as e:
|
||||
logger.warning(f"Backup scan error in {vault_backup_dir}: {e}")
|
||||
rows.sort(key=lambda r: r["timestamp"], reverse=True)
|
||||
return rows
|
||||
|
||||
|
||||
@router.get("/backup-stats")
|
||||
async def get_admin_backup_stats(_admin=Depends(require_admin)):
|
||||
"""Return backup statistics across all vaults."""
|
||||
rows = _scan_backups()
|
||||
now_ts = int(time.time())
|
||||
total_size_mb = round(sum(r["size"] for r in rows) / (1024 ** 2), 2)
|
||||
|
||||
by_vault: dict[str, dict[str, Any]] = {}
|
||||
for r in rows:
|
||||
bucket = by_vault.setdefault(r["vault"], {"count": 0, "size_mb": 0.0})
|
||||
bucket["count"] += 1
|
||||
bucket["size_mb"] = round(bucket["size_mb"] + r["size"] / (1024 ** 2), 2)
|
||||
|
||||
if rows:
|
||||
newest_ts = rows[0]["timestamp"]
|
||||
oldest_ts = rows[-1]["timestamp"]
|
||||
newest_age_days = round((now_ts - newest_ts) / 86400, 2)
|
||||
oldest_age_days = round((now_ts - oldest_ts) / 86400, 2)
|
||||
else:
|
||||
newest_age_days = 0.0
|
||||
oldest_age_days = 0.0
|
||||
|
||||
return {
|
||||
"total_backups": len(rows),
|
||||
"total_size_mb": total_size_mb,
|
||||
"oldest_age_days": oldest_age_days,
|
||||
"newest_age_days": newest_age_days,
|
||||
"by_vault": by_vault,
|
||||
}
|
||||
|
||||
|
||||
# ── /api/admin/stream — Server-Sent Events ─────────────────────────────
|
||||
async def _stats_event_generator():
|
||||
"""Yield an SSE `stats` event with the current system stats."""
|
||||
try:
|
||||
import psutil
|
||||
# Prime cpu_percent so the first real measurement isn't 0.0
|
||||
psutil.cpu_percent(interval=None)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
last_keepalive = time.time()
|
||||
KEEPALIVE_INTERVAL = 15.0
|
||||
STATS_INTERVAL = 5.0
|
||||
|
||||
try:
|
||||
while True:
|
||||
now = time.time()
|
||||
try:
|
||||
data = await asyncio.to_thread(_build_stats_payload)
|
||||
except Exception as e:
|
||||
logger.warning(f"SSE stats build failed: {e}")
|
||||
data = {"error": str(e)}
|
||||
|
||||
yield f"event: stats\ndata: {_format_timestamp_for_stream(data)}\n\n"
|
||||
|
||||
if now - last_keepalive >= KEEPALIVE_INTERVAL:
|
||||
last_keepalive = now
|
||||
yield ": keepalive\n\n"
|
||||
|
||||
await asyncio.sleep(STATS_INTERVAL)
|
||||
except asyncio.CancelledError:
|
||||
logger.info("SSE stats stream cancelled")
|
||||
raise
|
||||
|
||||
|
||||
def _build_stats_payload() -> dict:
|
||||
"""Synchronous stats builder for the SSE generator."""
|
||||
import psutil
|
||||
cpu_pct = psutil.cpu_percent(interval=None)
|
||||
vm = psutil.virtual_memory()
|
||||
disk_used_gb, disk_total_gb = _get_disk_stats()
|
||||
return {
|
||||
"cpu_pct": cpu_pct,
|
||||
"mem_used_mb": round(vm.used / (1024 ** 2), 1),
|
||||
"mem_total_mb": round(vm.total / (1024 ** 2), 1),
|
||||
"disk_used_gb": disk_used_gb,
|
||||
"disk_total_gb": disk_total_gb,
|
||||
"uptime_seconds": int(_server_uptime_seconds()),
|
||||
"timestamp": datetime.now(timezone.utc).isoformat(),
|
||||
}
|
||||
|
||||
|
||||
@router.get("/stream")
|
||||
async def stream_admin_stats(_admin=Depends(require_admin)):
|
||||
"""Server-Sent Events stream of system metrics, every 5 seconds."""
|
||||
return StreamingResponse(
|
||||
_stats_event_generator(),
|
||||
media_type="text/event-stream",
|
||||
headers={
|
||||
"Cache-Control": "no-cache",
|
||||
"X-Accel-Buffering": "no",
|
||||
"Connection": "keep-alive",
|
||||
},
|
||||
)
|
||||
+13
-5
@@ -662,11 +662,11 @@ class SSESafeGZipMiddleware(GZipMiddleware):
|
||||
We detect SSE endpoints by path and bypass compression entirely.
|
||||
"""
|
||||
async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None:
|
||||
if scope["type"] == "http" and scope.get("path") == "/api/events":
|
||||
# Bypass GZip: passthrough directly to the inner app
|
||||
await self.app(scope, receive, send)
|
||||
else:
|
||||
await super().__call__(scope, receive, send)
|
||||
if scope["type"] == "http" and scope.get("path") in ("/api/events", "/api/admin/stream"):
|
||||
# Bypass GZip: passthrough directly to the inner app
|
||||
await self.app(scope, receive, send)
|
||||
else:
|
||||
await super().__call__(scope, receive, send)
|
||||
|
||||
app.add_middleware(SSESafeGZipMiddleware, minimum_size=1000)
|
||||
|
||||
@@ -717,6 +717,14 @@ app.include_router(auth_router)
|
||||
app.include_router(ai_router)
|
||||
app.include_router(bookslm_router)
|
||||
|
||||
# Admin Dashboard endpoints (system stats, audit logs, backups, stream)
|
||||
try:
|
||||
from backend.admin import router as admin_router
|
||||
app.include_router(admin_router)
|
||||
logger.info("Admin dashboard router mounted at /api/admin/*")
|
||||
except ImportError as e:
|
||||
logger.warning(f"Could not load admin dashboard router: {e}")
|
||||
|
||||
# Resolve frontend path relative to this file
|
||||
FRONTEND_DIR = Path(__file__).resolve().parent.parent / "frontend"
|
||||
|
||||
|
||||
@@ -14,3 +14,4 @@ weasyprint>=60.0
|
||||
httpx>=0.27.0
|
||||
pypdf>=4.0
|
||||
pyotp>=2.10.0
|
||||
psutil>=5.9
|
||||
|
||||
Reference in New Issue
Block a user