"""Private management API (/api/*): auth, models, system, runtime, downloads, harvest, keys, settings, events.""" from __future__ import annotations import asyncio import json import time from pathlib import Path from typing import Any from fastapi import APIRouter, Depends, Query, Request, Response from fastapi.responses import StreamingResponse from pydantic import BaseModel, Field from ..auth import SESSION_COOKIE, Principal, get_auth, login_limiter, require_admin from ..errors import APIError, AuthError, Conflict, ModelNotFound, RateLimited from ..events import bus from ..hardware import GB, detect_hardware router = APIRouter(prefix="/api") SETTINGS_KEYS = { "max_model_memory_gb": float, "absolute_max_memory_gb": float, "min_free_disk_gb": float, "model_idle_timeout_minutes": int, "max_simultaneous_models": int, "preferred_runtime": str, "default_model": str, "preload_model": str, "log_prompts": bool, "allow_downloads": bool, "allow_gguf": bool, "allow_mlx": bool, "default_context": int, "default_max_tokens": int, "benchmark_max_tokens": int, "benchmark_runs": int, } def _ip(request: Request) -> str: fwd = request.headers.get("x-forwarded-for") return fwd.split(",")[0].strip() if fwd else (request.client.host if request.client else "?") def _actor(request: Request) -> str: p = getattr(request.state, "principal", None) return f"{p.kind}:{p.name}" if p else "anonymous" # --------------------------------------------------------------------------- # Auth # --------------------------------------------------------------------------- class LoginBody(BaseModel): email: str password: str class SetupBody(BaseModel): email: str password: str = Field(min_length=10) class PasswordBody(BaseModel): current_password: str new_password: str = Field(min_length=10) def _set_cookie(response: Response, request: Request, token: str) -> None: settings = request.app.state.settings secure = settings.secure_cookies or request.url.scheme == "https" or request.headers.get("x-forwarded-proto") == "https" response.set_cookie(SESSION_COOKIE, token, httponly=True, secure=secure, samesite="lax", max_age=settings.session_hours * 3600, path="/") @router.get("/auth/status") async def auth_status(request: Request): auth = get_auth(request) from ..auth import principal_from_request p = await principal_from_request(request, auth) return {"needs_setup": (await auth.user_count()) == 0, "authenticated": bool(p and p.has("admin")), "principal": {"kind": p.kind, "name": p.name} if p else None} @router.post("/auth/setup") async def auth_setup(body: SetupBody, request: Request, response: Response): auth = get_auth(request) if await auth.user_count() > 0: raise Conflict("Setup already completed.") uid = await auth.create_user(body.email, body.password) token = auth.make_session(uid, body.email.lower()) _set_cookie(response, request, token) await request.app.state.db.audit("auth.setup", actor=body.email, ip=_ip(request)) return {"ok": True, "email": body.email.lower()} @router.post("/auth/login") async def auth_login(body: LoginBody, request: Request, response: Response): auth = get_auth(request) if not login_limiter.check(_ip(request)): raise RateLimited("Too many login attempts. Try again in a minute.") user = await auth.authenticate(body.email, body.password) if not user: await request.app.state.db.audit("auth.login_failed", actor=body.email, ip=_ip(request)) raise AuthError("Invalid email or password.", code="INVALID_CREDENTIALS") token = auth.make_session(user["id"], user["email"]) _set_cookie(response, request, token) await request.app.state.db.audit("auth.login", actor=user["email"], ip=_ip(request)) return {"ok": True, "email": user["email"]} @router.post("/auth/logout") async def auth_logout(request: Request, response: Response): response.delete_cookie(SESSION_COOKIE, path="/") return {"ok": True} @router.get("/auth/me") async def auth_me(request: Request, p: Principal = Depends(require_admin)): return {"kind": p.kind, "name": p.name, "scopes": sorted(p.scopes)} @router.post("/auth/password") async def auth_password(body: PasswordBody, request: Request, p: Principal = Depends(require_admin)): auth = get_auth(request) if p.kind != "session": raise APIError("Password changes require a dashboard session.") user = await auth.authenticate(p.name, body.current_password) if not user: raise AuthError("Current password is incorrect.", code="INVALID_CREDENTIALS") await auth.change_password(user["id"], body.new_password) await request.app.state.db.audit("auth.password_changed", actor=p.name, ip=_ip(request)) return {"ok": True} # --------------------------------------------------------------------------- # Models # --------------------------------------------------------------------------- def _decorate(manager, m: dict) -> dict: st = manager.status_of(m["id"]) m["status"] = st m["loaded"] = st == "ready" lm = manager.loaded.get(m["id"]) if lm: m["worker"] = lm.to_dict() p = manager.progress.get(m["id"]) if p: m["progress"] = p return m @router.get("/models") async def list_models(request: Request, p: Principal = Depends(require_admin), include_missing: bool = False): manager = request.app.state.manager models = await manager.registry.list_models(include_missing=include_missing) aliases = await manager.registry.aliases() by_model: dict[str, list[str]] = {} for a, mid in aliases.items(): by_model.setdefault(mid, []).append(a) for m in models: _decorate(manager, m) m["aliases"] = by_model.get(m["id"], []) return {"models": models, "aliases": aliases, "count": len(models)} @router.get("/models/{model_id}") async def get_model(model_id: str, request: Request, p: Principal = Depends(require_admin)): manager = request.app.state.manager m = await manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") _decorate(manager, m) db = request.app.state.db m["benchmarks"] = await db.fetchall("SELECT * FROM model_benchmarks WHERE model_id=? ORDER BY created_at DESC LIMIT 30", (m["id"],)) for b in m["benchmarks"]: for k in ("params", "notes"): if b.get(k): try: b[k] = json.loads(b[k]) except Exception: pass m["events"] = await db.fetchall("SELECT * FROM model_events WHERE model_id=? ORDER BY created_at DESC LIMIT 50", (m["id"],)) m["aliases"] = [a for a, mid in (await manager.registry.aliases()).items() if mid == m["id"]] m["recent_requests"] = await db.fetchall( "SELECT created_at, endpoint, prompt_tokens, completion_tokens, ttft_ms, total_ms, tps, status FROM inference_requests " "WHERE model_id=? ORDER BY created_at DESC LIMIT 20", (m["id"],)) from ..models.estimator import CONTEXT_STEPS, estimate rt = m["runtime"] m["memory_curve"] = [estimate(m["weights_bytes"], rt, m["kv_bytes_per_token"] or 0, c, m["vision"]).to_dict() for c in CONTEXT_STEPS if not m["max_context"] or c <= max(m["max_context"], 2048)] return m @router.get("/models/{model_id}/files") async def model_files(model_id: str, request: Request, p: Principal = Depends(require_admin)): manager = request.app.state.manager m = await manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") return {"path": m["path"], "files": await manager.registry.files(m["id"])} class LoadBody(BaseModel): context: int | None = None force: bool = False @router.post("/models/{model_id}/load") async def load_model(model_id: str, request: Request, body: LoadBody | None = None, p: Principal = Depends(require_admin)): manager = request.app.state.manager body = body or LoadBody() lm = await manager.load(model_id, context=body.context, force=body.force) await request.app.state.db.audit("model.load", actor=_actor(request), target=lm.model["id"], ip=_ip(request)) return {"ok": True, "worker": lm.to_dict()} @router.post("/models/{model_id}/unload") async def unload_model(model_id: str, request: Request, p: Principal = Depends(require_admin)): manager = request.app.state.manager m = await manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") ok = await manager.unload(m["id"], reason="manual") await request.app.state.db.audit("model.unload", actor=_actor(request), target=m["id"], ip=_ip(request)) return {"ok": ok, "status": manager.status_of(m["id"])} class PatchModel(BaseModel): favorite: bool | None = None pinned: bool | None = None enabled: bool | None = None notes: str | None = None tags: list[str] | None = None name: str | None = None overrides: dict[str, Any] | None = None @router.patch("/models/{model_id}") async def patch_model(model_id: str, body: PatchModel, request: Request, p: Principal = Depends(require_admin)): manager = request.app.state.manager m = await manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") fields = {k: v for k, v in body.model_dump().items() if v is not None} if "overrides" in fields: ov = fields["overrides"] allowed = {"task", "vision", "embedding", "reranker", "family", "quantization", "max_context", "context", "kv_bits", "parameter_count", "pooling", "llama_args", "thinking", "tools", "name"} fields["overrides"] = {k: v for k, v in ov.items() if k in allowed} if fields: await manager.registry.update(m["id"], **fields) if "overrides" in fields: await manager.registry.rescan() await manager.refresh_model(m["id"]) await request.app.state.db.audit("model.update", actor=_actor(request), target=m["id"], detail=fields, ip=_ip(request)) return _decorate(manager, await manager.registry.get(m["id"])) # type: ignore[arg-type] @router.post("/models/{model_id}/pin") async def pin_model(model_id: str, request: Request, p: Principal = Depends(require_admin), pinned: bool = True): manager = request.app.state.manager m = await manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") await manager.registry.update(m["id"], pinned=pinned) await manager.refresh_model(m["id"]) return {"ok": True, "pinned": pinned} @router.post("/models/{model_id}/favorite") async def favorite_model(model_id: str, request: Request, p: Principal = Depends(require_admin), favorite: bool = True): manager = request.app.state.manager m = await manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") await manager.registry.update(m["id"], favorite=favorite) return {"ok": True, "favorite": favorite} class DeleteBody(BaseModel): confirm: str keep_benchmarks: bool = True @router.delete("/models/{model_id}") async def delete_model(model_id: str, body: DeleteBody, request: Request, p: Principal = Depends(require_admin)): manager = request.app.state.manager m = await manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") if body.confirm != m["id"]: raise APIError("Type the exact model id in 'confirm' to delete it.", code="CONFIRMATION_REQUIRED") if manager.loaded.get(m["id"]): raise Conflict("Unload the model before deleting it.") return await request.app.state.downloader.delete_model(m["id"], actor=_actor(request), keep_benchmarks=body.keep_benchmarks) @router.post("/models/rescan") async def rescan(request: Request, p: Principal = Depends(require_admin)): state = request.app.state async def run(job): def prog(frac, name): state.jobs.update(job, progress=frac, current=name) summary = await state.manager.registry.rescan(progress=prog) bus.publish("models", {"event": "rescan", **summary}) return summary job = state.jobs.submit("scan", "Rescan model directory", {}, run) # Scans are quick: wait a bit so callers get the result directly when possible for _ in range(100): await asyncio.sleep(0.1) if job.status in ("completed", "failed"): break return {"job": job.to_dict()} class BenchBody(BaseModel): max_tokens: int = 256 runs: int = 2 long_prompt: bool = True @router.post("/models/{model_id}/benchmark") async def benchmark(model_id: str, request: Request, body: BenchBody | None = None, p: Principal = Depends(require_admin)): state = request.app.state m = await state.manager.registry.resolve(model_id) if not m: raise ModelNotFound(f"Model '{model_id}' not found.") body = body or BenchBody() from ..bench import run_benchmark async def run(job): return await run_benchmark(state, m["id"], job, body.model_dump()) job = state.jobs.submit("benchmark", f"Benchmark {m['name']}", {"model_id": m["id"], **body.model_dump()}, run) return {"job": job.to_dict()} @router.get("/models/{model_id}/benchmarks") async def benchmarks(model_id: str, request: Request, p: Principal = Depends(require_admin)): rows = await request.app.state.db.fetchall("SELECT * FROM model_benchmarks WHERE model_id=? ORDER BY created_at DESC LIMIT 100", (model_id,)) for b in rows: for k in ("params", "notes"): if b.get(k): try: b[k] = json.loads(b[k]) except Exception: pass return {"benchmarks": rows} # --------------------------------------------------------------------------- # Aliases # --------------------------------------------------------------------------- class AliasBody(BaseModel): alias: str = Field(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9][A-Za-z0-9._-]*$") model_id: str @router.get("/aliases") async def list_aliases(request: Request, p: Principal = Depends(require_admin)): return {"aliases": await request.app.state.manager.registry.aliases()} @router.put("/aliases") async def put_alias(body: AliasBody, request: Request, p: Principal = Depends(require_admin)): reg = request.app.state.manager.registry m = await reg.get(body.model_id) if not m: raise ModelNotFound(f"Model '{body.model_id}' not found.") if await reg.get(body.alias): raise Conflict("An installed model already has this id.") if body.alias == "auto": raise APIError("'auto' is reserved.") await reg.set_alias(body.alias, m["id"]) await request.app.state.db.audit("alias.set", actor=_actor(request), target=body.alias, detail=body.model_id) return {"aliases": await reg.aliases()} @router.delete("/aliases/{alias}") async def delete_alias(alias: str, request: Request, p: Principal = Depends(require_admin)): reg = request.app.state.manager.registry await reg.delete_alias(alias) return {"aliases": await reg.aliases()} # --------------------------------------------------------------------------- # System / runtime # --------------------------------------------------------------------------- @router.get("/system") async def system(request: Request, p: Principal = Depends(require_admin)): state = request.app.state hw = detect_hardware(state.settings.models_dir).to_dict() tel = state.metrics.last or await state.metrics.sample() budget, absolute, max_models = await state.manager.budgets() return {"hardware": hw, "telemetry": tel, "manager": state.manager.snapshot(), "policy": {"max_model_memory_gb": budget, "absolute_max_memory_gb": absolute, "max_simultaneous_models": max_models, "macos_reserve_gb": state.settings.macos_reserve_gb, "min_free_disk_gb": float(await state.db.get_setting("min_free_disk_gb", state.settings.min_free_disk_gb))}, "version": request.app.version, "started_at": state.metrics.started_at, "runtimes": await _runtime_versions()} async def _runtime_versions() -> dict: out: dict[str, Any] = {} try: import mlx.core as mx import mlx_lm out["mlx"] = mx.__version__ out["mlx_lm"] = mlx_lm.__version__ except Exception: out["mlx"] = None try: import mlx_vlm out["mlx_vlm"] = mlx_vlm.__version__ except Exception: out["mlx_vlm"] = None try: import shutil import subprocess b = shutil.which("llama-server") or "/opt/homebrew/bin/llama-server" r = subprocess.run([b, "--version"], capture_output=True, text=True, timeout=5) out["llama_cpp"] = (r.stdout + r.stderr).strip().splitlines()[0][:80] if (r.stdout or r.stderr) else None except Exception: out["llama_cpp"] = None return out @router.get("/system/memory") async def system_memory(request: Request, p: Principal = Depends(require_admin)): tel = await request.app.state.metrics.sample() return {k: tel[k] for k in tel if k.startswith(("mem_", "swap_", "worker_", "loaded_"))} @router.get("/system/gpu") async def system_gpu(request: Request, p: Principal = Depends(require_admin)): tel = await request.app.state.metrics.sample() hw = detect_hardware(request.app.state.settings.models_dir) return {"gpu_cores": hw.gpu_cores, "chip": hw.chip, "gpu_percent": tel["gpu_percent"], "gpu_renderer_percent": tel["gpu_renderer_percent"], "gpu_memory_gb": tel["gpu_memory_gb"], "thermal_state": tel["thermal_state"]} @router.get("/system/storage") async def system_storage(request: Request, p: Principal = Depends(require_admin)): state = request.app.state settings = state.settings import shutil def du(path: Path) -> int: total = 0 if path.exists(): for f in path.rglob("*"): if f.is_file(): try: total += f.stat().st_size except OSError: pass return total usage = shutil.disk_usage(settings.models_dir) models_bytes, logs_bytes, db_bytes = await asyncio.gather( asyncio.to_thread(du, settings.models_dir), asyncio.to_thread(du, settings.logs_path), asyncio.to_thread(lambda: sum(du(p) for p in [settings.db_path] if p.exists()) + du(settings.data_path))) hf_cache = Path.home() / ".cache" / "huggingface" cache_bytes = await asyncio.to_thread(du, hf_cache) by_model = await state.db.fetchall("SELECT id, name, disk_size_bytes, runtime, last_used_at FROM models WHERE installed=1 ORDER BY disk_size_bytes DESC") return {"total_gb": round(usage.total / GB, 1), "used_gb": round(usage.used / GB, 1), "free_gb": round(usage.free / GB, 1), "models_gb": round(models_bytes / GB, 2), "logs_gb": round(logs_bytes / GB, 3), "database_gb": round(db_bytes / GB, 3), "cache_gb": round(cache_bytes / GB, 2), "other_gb": round(max(0, usage.used - models_bytes - logs_bytes - db_bytes - cache_bytes) / GB, 1), "model_root": str(settings.models_dir), "models": by_model, "min_free_gb": float(await state.db.get_setting("min_free_disk_gb", settings.min_free_disk_gb))} @router.get("/system/processes") async def system_processes(request: Request, p: Principal = Depends(require_admin)): import psutil manager = request.app.state.manager procs = [] for lm in manager.loaded.values(): try: pr = psutil.Process(lm.handle.pid) procs.append({"model_id": lm.model["id"], "pid": lm.handle.pid, "port": lm.handle.port, "runtime": lm.handle.runtime, "rss_gb": round(lm.handle.memory_bytes() / GB, 2), "cpu_percent": pr.cpu_percent(interval=None), "status": lm.status, "threads": pr.num_threads(), "created": pr.create_time()}) except psutil.Error: pass me = psutil.Process() return {"workers": procs, "server": {"pid": me.pid, "rss_gb": round(me.memory_info().rss / GB, 3), "cpu_percent": me.cpu_percent(interval=None), "threads": me.num_threads()}} @router.get("/system/health") async def system_health_admin(request: Request): return await public_health(request) async def public_health(request: Request) -> dict: state = request.app.state hw = detect_hardware(state.settings.models_dir) tel = state.metrics.last or await state.metrics.sample() cur = state.manager.current_model() return {"status": "ok", "hardware": {"chip": hw.chip, "memory_gb": hw.memory_gb, "gpu_cores": hw.gpu_cores}, "model": {"loaded": cur is not None, "id": cur.model["id"] if cur else None, "loaded_models": [lm.model["id"] for lm in state.manager.loaded.values() if lm.status == "ready"]}, "memory": {"used_gb": tel["mem_used_gb"], "available_gb": tel["mem_available_gb"], "pressure": tel["mem_pressure_level"]}, "uptime_seconds": round(time.time() - state.metrics.started_at), "version": request.app.version} @router.get("/system/metrics") async def system_metrics(request: Request, p: Principal = Depends(require_admin), minutes: int = Query(60, ge=5, le=60 * 24 * 7), hours: int = Query(24, ge=1, le=24 * 30)): state = request.app.state return {"history": await state.metrics.history(minutes), "requests": await state.metrics.request_stats(hours), "live": state.metrics.last} @router.get("/runtime/status") async def runtime_status(request: Request, p: Principal = Depends(require_admin)): return request.app.state.manager.snapshot() @router.get("/runtime/current-model") async def runtime_current(request: Request, p: Principal = Depends(require_admin)): cur = request.app.state.manager.current_model() return {"model": cur.to_dict() if cur else None} # --------------------------------------------------------------------------- # Jobs / downloads / harvest # --------------------------------------------------------------------------- @router.get("/jobs") async def jobs(request: Request, p: Principal = Depends(require_admin), kind: str | None = None, limit: int = 100): kinds = {kind} if kind else None return {"jobs": await request.app.state.jobs.history(kinds, limit)} @router.post("/jobs/{job_id}/cancel") async def cancel_job(job_id: str, request: Request, p: Principal = Depends(require_admin)): ok = await request.app.state.jobs.cancel(job_id) return {"ok": ok} class InspectBody(BaseModel): repository: str quant: str | None = None class DownloadBody(BaseModel): repository: str quant: str | None = None force: bool = False @router.post("/models/inspect") async def inspect_repo(body: InspectBody, request: Request, p: Principal = Depends(require_admin)): return await request.app.state.downloader.inspect(body.repository, body.quant) @router.post("/models/download") async def download(body: DownloadBody, request: Request, p: Principal = Depends(require_admin)): job = await request.app.state.downloader.start_download(body.repository, body.quant, force=body.force, actor=_actor(request)) return {"job": job.to_dict()} @router.get("/downloads") async def downloads(request: Request, p: Principal = Depends(require_admin)): return {"downloads": await request.app.state.jobs.history({"download"}, 100)} @router.post("/downloads/{job_id}/retry") async def retry_download(job_id: str, request: Request, p: Principal = Depends(require_admin)): rows = await request.app.state.db.fetchone("SELECT payload FROM jobs WHERE id=?", (job_id,)) if not rows: raise APIError("Job not found.", status_code=404, code="JOB_NOT_FOUND") payload = json.loads(rows["payload"] or "{}") job = await request.app.state.downloader.start_download(payload["repository"], payload.get("quant"), force=True, actor=_actor(request)) return {"job": job.to_dict()} class HarvestBody(BaseModel): runtimes: list[str] | None = None authors: list[str] | None = None limit_per_author: int = 150 min_downloads: int = 500 max_ram_gb: float | None = None families: list[str] | None = None tasks: list[str] | None = None search: str | None = None @router.post("/harvest/scan") async def harvest_scan(body: HarvestBody, request: Request, p: Principal = Depends(require_admin)): job = await request.app.state.harvester.start_scan(body.model_dump(), actor=_actor(request)) return {"job": job.to_dict()} @router.get("/harvest/candidates") async def harvest_candidates(request: Request, p: Principal = Depends(require_admin), include_duplicates: bool = False, task: str | None = None, runtime: str | None = None, family: str | None = None, size_class: str | None = None, q: str | None = None, limit: int = 300): h = request.app.state.harvester rows = await h.candidates(include_duplicates=include_duplicates, task=task, runtime=runtime, family=family, size_class=size_class, q=q, limit=limit) last = await request.app.state.db.scalar("SELECT MAX(scanned_at) FROM harvest_candidates") return {"candidates": rows, "last_scan": last, "starter": await h.suggest_starter()} class SelectBody(BaseModel): repo_id: str selected: bool = True @router.post("/harvest/select") async def harvest_select(body: SelectBody, request: Request, p: Principal = Depends(require_admin)): await request.app.state.harvester.select(body.repo_id, body.selected) return {"ok": True} @router.post("/harvest/dismiss") async def harvest_dismiss(body: SelectBody, request: Request, p: Principal = Depends(require_admin)): await request.app.state.harvester.dismiss(body.repo_id) return {"ok": True} @router.post("/harvest/queue") async def harvest_queue(request: Request, p: Principal = Depends(require_admin)): return {"queued": await request.app.state.harvester.queue_selected(actor=_actor(request))} # --------------------------------------------------------------------------- # API keys # --------------------------------------------------------------------------- class KeyBody(BaseModel): name: str = Field(min_length=1, max_length=64) scopes: list[str] = ["inference"] @router.get("/keys") async def list_keys(request: Request, p: Principal = Depends(require_admin)): rows = await request.app.state.db.fetchall( "SELECT id, name, prefix, scopes, created_at, last_used_at, request_count, revoked_at FROM api_keys ORDER BY created_at DESC") for r in rows: r["scopes"] = r["scopes"].split(",") return {"keys": rows} @router.post("/keys") async def create_key(body: KeyBody, request: Request, p: Principal = Depends(require_admin)): scopes = [s for s in body.scopes if s in ("inference", "admin")] or ["inference"] raw, row = await get_auth(request).create_key(body.name, scopes) await request.app.state.db.audit("key.create", actor=_actor(request), target=body.name, ip=_ip(request)) row = dict(row) row["scopes"] = row["scopes"].split(",") return {"key": raw, "record": row} class KeyPatch(BaseModel): name: str | None = None @router.patch("/keys/{key_id}") async def rename_key(key_id: int, body: KeyPatch, request: Request, p: Principal = Depends(require_admin)): if body.name: await request.app.state.db.execute("UPDATE api_keys SET name=? WHERE id=?", (body.name, key_id)) return {"ok": True} @router.delete("/keys/{key_id}") async def revoke_key(key_id: int, request: Request, p: Principal = Depends(require_admin)): await request.app.state.db.execute("UPDATE api_keys SET revoked_at=? WHERE id=? AND revoked_at IS NULL", (time.time(), key_id)) await request.app.state.db.audit("key.revoke", actor=_actor(request), target=str(key_id), ip=_ip(request)) return {"ok": True} # --------------------------------------------------------------------------- # Settings / logs / events # --------------------------------------------------------------------------- @router.get("/settings") async def get_settings_(request: Request, p: Principal = Depends(require_admin)): state = request.app.state s = state.settings defaults = { "max_model_memory_gb": s.max_model_memory_gb, "absolute_max_memory_gb": s.absolute_max_memory_gb, "min_free_disk_gb": s.min_free_disk_gb, "model_idle_timeout_minutes": s.model_idle_timeout_minutes, "max_simultaneous_models": s.max_simultaneous_models, "preferred_runtime": "mlx", "default_model": None, "preload_model": s.preload_model, "log_prompts": s.log_prompts, "allow_downloads": s.allow_downloads, "allow_gguf": s.enable_gguf, "allow_mlx": s.enable_mlx, "default_context": s.default_context, "default_max_tokens": s.default_max_tokens, "benchmark_max_tokens": 256, "benchmark_runs": 2, } stored = await state.db.all_settings() merged = {**defaults, **{k: v for k, v in stored.items() if k in defaults}} return {"settings": merged, "defaults": defaults, "paths": {"root": str(s.root), "models": str(s.models_dir), "data": str(s.data_path), "logs": str(s.logs_path), "db": str(s.db_path)}, "hf_token_set": bool(s.hf_token), "public_url": s.public_url} @router.patch("/settings") async def patch_settings(request: Request, p: Principal = Depends(require_admin)): state = request.app.state body = await request.json() if not isinstance(body, dict): raise APIError("Body must be an object.") changed = {} hw = detect_hardware(state.settings.models_dir) for k, v in body.items(): if k not in SETTINGS_KEYS: continue typ = SETTINGS_KEYS[k] try: if typ is bool: v = bool(v) elif v is None or v == "": v = None else: v = typ(v) except (TypeError, ValueError): raise APIError(f"Invalid value for {k}.", param=k) if k in ("max_model_memory_gb", "absolute_max_memory_gb") and v is not None: if v < 1 or v > hw.memory_gb - 4: raise APIError(f"{k} must be between 1 and {hw.memory_gb - 4:.0f} GB on this machine.", param=k) if k == "max_simultaneous_models" and v is not None and (v < 1 or v > 4): raise APIError("max_simultaneous_models must be 1-4.", param=k) changed[k] = v await state.db.set_setting(k, v) if changed: await state.db.audit("settings.update", actor=_actor(request), detail=changed, ip=_ip(request)) if any(k in changed for k in ("max_model_memory_gb", "absolute_max_memory_gb")): await state.manager.registry.reevaluate_all() bus.publish("settings", changed) return {"changed": changed} @router.get("/logs/audit") async def audit_logs(request: Request, p: Principal = Depends(require_admin), limit: int = 200): rows = await request.app.state.db.fetchall("SELECT * FROM audit_logs ORDER BY created_at DESC LIMIT ?", (limit,)) return {"logs": rows} @router.get("/logs/events") async def model_events(request: Request, p: Principal = Depends(require_admin), limit: int = 200): rows = await request.app.state.db.fetchall("SELECT * FROM model_events ORDER BY created_at DESC LIMIT ?", (limit,)) return {"events": rows} @router.get("/logs/requests") async def request_logs(request: Request, p: Principal = Depends(require_admin), limit: int = 200): rows = await request.app.state.db.fetchall( "SELECT id, created_at, model_id, requested_model, endpoint, api_key_id, stream, prompt_tokens, completion_tokens, ttft_ms, " "total_ms, tps, load_wait_ms, status, error_code FROM inference_requests ORDER BY created_at DESC LIMIT ?", (limit,)) return {"requests": rows} @router.get("/logs/worker/{model_id}") async def worker_log(model_id: str, request: Request, p: Principal = Depends(require_admin), lines: int = 200): path = request.app.state.settings.logs_path / "workers" / f"worker-{model_id}.log" if not path.exists(): return {"lines": []} txt = path.read_text(errors="replace").splitlines() return {"lines": txt[-lines:]} @router.get("/events") async def events(request: Request, p: Principal = Depends(require_admin)): q = bus.subscribe() state = request.app.state async def gen(): try: # initial snapshot snap = {"seq": 0, "ts": time.time(), "type": "snapshot", "data": {"manager": state.manager.snapshot(), "metrics": state.metrics.last, "jobs": state.jobs.list(limit=30)}} yield bus.format_sse(snap) while True: try: ev = await asyncio.wait_for(q.get(), timeout=15) yield bus.format_sse(ev) except asyncio.TimeoutError: yield ": ping\n\n" if await request.is_disconnected(): break finally: bus.unsubscribe(q) return StreamingResponse(gen(), media_type="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}) @router.post("/models/{model_id}/tokenize") async def tokenize(model_id: str, request: Request, p: Principal = Depends(require_admin)): manager = request.app.state.manager lm = manager.get_ready(model_id) if not lm: raise Conflict("Model is not loaded.") body = await request.json() if lm.handle.runtime == "mlx": r = await manager.client.post(f"{lm.handle.base_url}/tokenize", json=body, timeout=30) else: text = body.get("text") or " ".join(m.get("content", "") for m in body.get("messages", []) if isinstance(m.get("content"), str)) r = await manager.client.post(f"{lm.handle.base_url}/tokenize", json={"content": text}, timeout=30) d = r.json() return {"tokens": len(d.get("tokens", [])), "max_context": lm.handle.context} return r.json()