1"""Private management API (/api/*): auth, models, system, runtime, downloads, harvest, keys, settings, events."""23from __future__ import annotations45import asyncio6import json7import time8from pathlib import Path9from typing import Any1011from fastapi import APIRouter, Depends, Query, Request, Response12from fastapi.responses import StreamingResponse13from pydantic import BaseModel, Field1415from ..auth import SESSION_COOKIE, Principal, get_auth, login_limiter, require_admin16from ..errors import APIError, AuthError, Conflict, ModelNotFound, RateLimited17from ..events import bus18from ..hardware import GB, detect_hardware1920router = APIRouter(prefix="/api")2122SETTINGS_KEYS = {23 "max_model_memory_gb": float, "absolute_max_memory_gb": float, "min_free_disk_gb": float,24 "model_idle_timeout_minutes": int, "max_simultaneous_models": int, "preferred_runtime": str,25 "default_model": str, "preload_model": str, "log_prompts": bool, "allow_downloads": bool,26 "allow_gguf": bool, "allow_mlx": bool, "default_context": int, "default_max_tokens": int,27 "benchmark_max_tokens": int, "benchmark_runs": int,28}293031def _ip(request: Request) -> str:32 fwd = request.headers.get("x-forwarded-for")33 return fwd.split(",")[0].strip() if fwd else (request.client.host if request.client else "?")343536def _actor(request: Request) -> str:37 p = getattr(request.state, "principal", None)38 return f"{p.kind}:{p.name}" if p else "anonymous"394041# ---------------------------------------------------------------------------42# Auth43# ---------------------------------------------------------------------------444546class LoginBody(BaseModel):47 email: str48 password: str495051class SetupBody(BaseModel):52 email: str53 password: str = Field(min_length=10)545556class PasswordBody(BaseModel):57 current_password: str58 new_password: str = Field(min_length=10)596061def _set_cookie(response: Response, request: Request, token: str) -> None:62 settings = request.app.state.settings63 secure = settings.secure_cookies or request.url.scheme == "https" or request.headers.get("x-forwarded-proto") == "https"64 response.set_cookie(SESSION_COOKIE, token, httponly=True, secure=secure, samesite="lax",65 max_age=settings.session_hours * 3600, path="/")666768@router.get("/auth/status")69async def auth_status(request: Request):70 auth = get_auth(request)71 from ..auth import principal_from_request72 p = await principal_from_request(request, auth)73 return {"needs_setup": (await auth.user_count()) == 0, "authenticated": bool(p and p.has("admin")),74 "principal": {"kind": p.kind, "name": p.name} if p else None}757677@router.post("/auth/setup")78async def auth_setup(body: SetupBody, request: Request, response: Response):79 auth = get_auth(request)80 if await auth.user_count() > 0:81 raise Conflict("Setup already completed.")82 uid = await auth.create_user(body.email, body.password)83 token = auth.make_session(uid, body.email.lower())84 _set_cookie(response, request, token)85 await request.app.state.db.audit("auth.setup", actor=body.email, ip=_ip(request))86 return {"ok": True, "email": body.email.lower()}878889@router.post("/auth/login")90async def auth_login(body: LoginBody, request: Request, response: Response):91 auth = get_auth(request)92 if not login_limiter.check(_ip(request)):93 raise RateLimited("Too many login attempts. Try again in a minute.")94 user = await auth.authenticate(body.email, body.password)95 if not user:96 await request.app.state.db.audit("auth.login_failed", actor=body.email, ip=_ip(request))97 raise AuthError("Invalid email or password.", code="INVALID_CREDENTIALS")98 token = auth.make_session(user["id"], user["email"])99 _set_cookie(response, request, token)100 await request.app.state.db.audit("auth.login", actor=user["email"], ip=_ip(request))101 return {"ok": True, "email": user["email"]}102103104@router.post("/auth/logout")105async def auth_logout(request: Request, response: Response):106 response.delete_cookie(SESSION_COOKIE, path="/")107 return {"ok": True}108109110@router.get("/auth/me")111async def auth_me(request: Request, p: Principal = Depends(require_admin)):112 return {"kind": p.kind, "name": p.name, "scopes": sorted(p.scopes)}113114115@router.post("/auth/password")116async def auth_password(body: PasswordBody, request: Request, p: Principal = Depends(require_admin)):117 auth = get_auth(request)118 if p.kind != "session":119 raise APIError("Password changes require a dashboard session.")120 user = await auth.authenticate(p.name, body.current_password)121 if not user:122 raise AuthError("Current password is incorrect.", code="INVALID_CREDENTIALS")123 await auth.change_password(user["id"], body.new_password)124 await request.app.state.db.audit("auth.password_changed", actor=p.name, ip=_ip(request))125 return {"ok": True}126127128# ---------------------------------------------------------------------------129# Models130# ---------------------------------------------------------------------------131132133def _decorate(manager, m: dict) -> dict:134 st = manager.status_of(m["id"])135 m["status"] = st136 m["loaded"] = st == "ready"137 lm = manager.loaded.get(m["id"])138 if lm:139 m["worker"] = lm.to_dict()140 p = manager.progress.get(m["id"])141 if p:142 m["progress"] = p143 return m144145146@router.get("/models")147async def list_models(request: Request, p: Principal = Depends(require_admin), include_missing: bool = False):148 manager = request.app.state.manager149 models = await manager.registry.list_models(include_missing=include_missing)150 aliases = await manager.registry.aliases()151 by_model: dict[str, list[str]] = {}152 for a, mid in aliases.items():153 by_model.setdefault(mid, []).append(a)154 for m in models:155 _decorate(manager, m)156 m["aliases"] = by_model.get(m["id"], [])157 return {"models": models, "aliases": aliases, "count": len(models)}158159160@router.get("/models/{model_id}")161async def get_model(model_id: str, request: Request, p: Principal = Depends(require_admin)):162 manager = request.app.state.manager163 m = await manager.registry.resolve(model_id)164 if not m:165 raise ModelNotFound(f"Model '{model_id}' not found.")166 _decorate(manager, m)167 db = request.app.state.db168 m["benchmarks"] = await db.fetchall("SELECT * FROM model_benchmarks WHERE model_id=? ORDER BY created_at DESC LIMIT 30", (m["id"],))169 for b in m["benchmarks"]:170 for k in ("params", "notes"):171 if b.get(k):172 try:173 b[k] = json.loads(b[k])174 except Exception:175 pass176 m["events"] = await db.fetchall("SELECT * FROM model_events WHERE model_id=? ORDER BY created_at DESC LIMIT 50", (m["id"],))177 m["aliases"] = [a for a, mid in (await manager.registry.aliases()).items() if mid == m["id"]]178 m["recent_requests"] = await db.fetchall(179 "SELECT created_at, endpoint, prompt_tokens, completion_tokens, ttft_ms, total_ms, tps, status FROM inference_requests "180 "WHERE model_id=? ORDER BY created_at DESC LIMIT 20", (m["id"],))181 from ..models.estimator import CONTEXT_STEPS, estimate182 rt = m["runtime"]183 m["memory_curve"] = [estimate(m["weights_bytes"], rt, m["kv_bytes_per_token"] or 0, c, m["vision"]).to_dict()184 for c in CONTEXT_STEPS if not m["max_context"] or c <= max(m["max_context"], 2048)]185 return m186187188@router.get("/models/{model_id}/files")189async def model_files(model_id: str, request: Request, p: Principal = Depends(require_admin)):190 manager = request.app.state.manager191 m = await manager.registry.resolve(model_id)192 if not m:193 raise ModelNotFound(f"Model '{model_id}' not found.")194 return {"path": m["path"], "files": await manager.registry.files(m["id"])}195196197class LoadBody(BaseModel):198 context: int | None = None199 force: bool = False200201202@router.post("/models/{model_id}/load")203async def load_model(model_id: str, request: Request, body: LoadBody | None = None, p: Principal = Depends(require_admin)):204 manager = request.app.state.manager205 body = body or LoadBody()206 lm = await manager.load(model_id, context=body.context, force=body.force)207 await request.app.state.db.audit("model.load", actor=_actor(request), target=lm.model["id"], ip=_ip(request))208 return {"ok": True, "worker": lm.to_dict()}209210211@router.post("/models/{model_id}/unload")212async def unload_model(model_id: str, request: Request, p: Principal = Depends(require_admin)):213 manager = request.app.state.manager214 m = await manager.registry.resolve(model_id)215 if not m:216 raise ModelNotFound(f"Model '{model_id}' not found.")217 ok = await manager.unload(m["id"], reason="manual")218 await request.app.state.db.audit("model.unload", actor=_actor(request), target=m["id"], ip=_ip(request))219 return {"ok": ok, "status": manager.status_of(m["id"])}220221222class PatchModel(BaseModel):223 favorite: bool | None = None224 pinned: bool | None = None225 enabled: bool | None = None226 notes: str | None = None227 tags: list[str] | None = None228 name: str | None = None229 overrides: dict[str, Any] | None = None230231232@router.patch("/models/{model_id}")233async def patch_model(model_id: str, body: PatchModel, request: Request, p: Principal = Depends(require_admin)):234 manager = request.app.state.manager235 m = await manager.registry.resolve(model_id)236 if not m:237 raise ModelNotFound(f"Model '{model_id}' not found.")238 fields = {k: v for k, v in body.model_dump().items() if v is not None}239 if "overrides" in fields:240 ov = fields["overrides"]241 allowed = {"task", "vision", "embedding", "reranker", "family", "quantization", "max_context", "context", "kv_bits",242 "parameter_count", "pooling", "llama_args", "thinking", "tools", "name"}243 fields["overrides"] = {k: v for k, v in ov.items() if k in allowed}244 if fields:245 await manager.registry.update(m["id"], **fields)246 if "overrides" in fields:247 await manager.registry.rescan()248 await manager.refresh_model(m["id"])249 await request.app.state.db.audit("model.update", actor=_actor(request), target=m["id"], detail=fields, ip=_ip(request))250 return _decorate(manager, await manager.registry.get(m["id"])) # type: ignore[arg-type]251252253@router.post("/models/{model_id}/pin")254async def pin_model(model_id: str, request: Request, p: Principal = Depends(require_admin), pinned: bool = True):255 manager = request.app.state.manager256 m = await manager.registry.resolve(model_id)257 if not m:258 raise ModelNotFound(f"Model '{model_id}' not found.")259 await manager.registry.update(m["id"], pinned=pinned)260 await manager.refresh_model(m["id"])261 return {"ok": True, "pinned": pinned}262263264@router.post("/models/{model_id}/favorite")265async def favorite_model(model_id: str, request: Request, p: Principal = Depends(require_admin), favorite: bool = True):266 manager = request.app.state.manager267 m = await manager.registry.resolve(model_id)268 if not m:269 raise ModelNotFound(f"Model '{model_id}' not found.")270 await manager.registry.update(m["id"], favorite=favorite)271 return {"ok": True, "favorite": favorite}272273274class DeleteBody(BaseModel):275 confirm: str276 keep_benchmarks: bool = True277278279@router.delete("/models/{model_id}")280async def delete_model(model_id: str, body: DeleteBody, request: Request, p: Principal = Depends(require_admin)):281 manager = request.app.state.manager282 m = await manager.registry.resolve(model_id)283 if not m:284 raise ModelNotFound(f"Model '{model_id}' not found.")285 if body.confirm != m["id"]:286 raise APIError("Type the exact model id in 'confirm' to delete it.", code="CONFIRMATION_REQUIRED")287 if manager.loaded.get(m["id"]):288 raise Conflict("Unload the model before deleting it.")289 return await request.app.state.downloader.delete_model(m["id"], actor=_actor(request), keep_benchmarks=body.keep_benchmarks)290291292@router.post("/models/rescan")293async def rescan(request: Request, p: Principal = Depends(require_admin)):294 state = request.app.state295296 async def run(job):297 def prog(frac, name):298 state.jobs.update(job, progress=frac, current=name)299 summary = await state.manager.registry.rescan(progress=prog)300 bus.publish("models", {"event": "rescan", **summary})301 return summary302303 job = state.jobs.submit("scan", "Rescan model directory", {}, run)304 # Scans are quick: wait a bit so callers get the result directly when possible305 for _ in range(100):306 await asyncio.sleep(0.1)307 if job.status in ("completed", "failed"):308 break309 return {"job": job.to_dict()}310311312class BenchBody(BaseModel):313 max_tokens: int = 256314 runs: int = 2315 long_prompt: bool = True316317318@router.post("/models/{model_id}/benchmark")319async def benchmark(model_id: str, request: Request, body: BenchBody | None = None, p: Principal = Depends(require_admin)):320 state = request.app.state321 m = await state.manager.registry.resolve(model_id)322 if not m:323 raise ModelNotFound(f"Model '{model_id}' not found.")324 body = body or BenchBody()325 from ..bench import run_benchmark326327 async def run(job):328 return await run_benchmark(state, m["id"], job, body.model_dump())329330 job = state.jobs.submit("benchmark", f"Benchmark {m['name']}", {"model_id": m["id"], **body.model_dump()}, run)331 return {"job": job.to_dict()}332333334@router.get("/models/{model_id}/benchmarks")335async def benchmarks(model_id: str, request: Request, p: Principal = Depends(require_admin)):336 rows = await request.app.state.db.fetchall("SELECT * FROM model_benchmarks WHERE model_id=? ORDER BY created_at DESC LIMIT 100", (model_id,))337 for b in rows:338 for k in ("params", "notes"):339 if b.get(k):340 try:341 b[k] = json.loads(b[k])342 except Exception:343 pass344 return {"benchmarks": rows}345346347# ---------------------------------------------------------------------------348# Aliases349# ---------------------------------------------------------------------------350351352class AliasBody(BaseModel):353 alias: str = Field(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9][A-Za-z0-9._-]*$")354 model_id: str355356357@router.get("/aliases")358async def list_aliases(request: Request, p: Principal = Depends(require_admin)):359 return {"aliases": await request.app.state.manager.registry.aliases()}360361362@router.put("/aliases")363async def put_alias(body: AliasBody, request: Request, p: Principal = Depends(require_admin)):364 reg = request.app.state.manager.registry365 m = await reg.get(body.model_id)366 if not m:367 raise ModelNotFound(f"Model '{body.model_id}' not found.")368 if await reg.get(body.alias):369 raise Conflict("An installed model already has this id.")370 if body.alias == "auto":371 raise APIError("'auto' is reserved.")372 await reg.set_alias(body.alias, m["id"])373 await request.app.state.db.audit("alias.set", actor=_actor(request), target=body.alias, detail=body.model_id)374 return {"aliases": await reg.aliases()}375376377@router.delete("/aliases/{alias}")378async def delete_alias(alias: str, request: Request, p: Principal = Depends(require_admin)):379 reg = request.app.state.manager.registry380 await reg.delete_alias(alias)381 return {"aliases": await reg.aliases()}382383384# ---------------------------------------------------------------------------385# System / runtime386# ---------------------------------------------------------------------------387388389@router.get("/system")390async def system(request: Request, p: Principal = Depends(require_admin)):391 state = request.app.state392 hw = detect_hardware(state.settings.models_dir).to_dict()393 tel = state.metrics.last or await state.metrics.sample()394 budget, absolute, max_models = await state.manager.budgets()395 return {"hardware": hw, "telemetry": tel, "manager": state.manager.snapshot(),396 "policy": {"max_model_memory_gb": budget, "absolute_max_memory_gb": absolute, "max_simultaneous_models": max_models,397 "macos_reserve_gb": state.settings.macos_reserve_gb,398 "min_free_disk_gb": float(await state.db.get_setting("min_free_disk_gb", state.settings.min_free_disk_gb))},399 "version": request.app.version, "started_at": state.metrics.started_at,400 "runtimes": await _runtime_versions()}401402403async def _runtime_versions() -> dict:404 out: dict[str, Any] = {}405 try:406 import mlx.core as mx407 import mlx_lm408 out["mlx"] = mx.__version__409 out["mlx_lm"] = mlx_lm.__version__410 except Exception:411 out["mlx"] = None412 try:413 import mlx_vlm414 out["mlx_vlm"] = mlx_vlm.__version__415 except Exception:416 out["mlx_vlm"] = None417 try:418 import shutil419 import subprocess420 b = shutil.which("llama-server") or "/opt/homebrew/bin/llama-server"421 r = subprocess.run([b, "--version"], capture_output=True, text=True, timeout=5)422 out["llama_cpp"] = (r.stdout + r.stderr).strip().splitlines()[0][:80] if (r.stdout or r.stderr) else None423 except Exception:424 out["llama_cpp"] = None425 return out426427428@router.get("/system/memory")429async def system_memory(request: Request, p: Principal = Depends(require_admin)):430 tel = await request.app.state.metrics.sample()431 return {k: tel[k] for k in tel if k.startswith(("mem_", "swap_", "worker_", "loaded_"))}432433434@router.get("/system/gpu")435async def system_gpu(request: Request, p: Principal = Depends(require_admin)):436 tel = await request.app.state.metrics.sample()437 hw = detect_hardware(request.app.state.settings.models_dir)438 return {"gpu_cores": hw.gpu_cores, "chip": hw.chip, "gpu_percent": tel["gpu_percent"],439 "gpu_renderer_percent": tel["gpu_renderer_percent"], "gpu_memory_gb": tel["gpu_memory_gb"],440 "thermal_state": tel["thermal_state"]}441442443@router.get("/system/storage")444async def system_storage(request: Request, p: Principal = Depends(require_admin)):445 state = request.app.state446 settings = state.settings447 import shutil448449 def du(path: Path) -> int:450 total = 0451 if path.exists():452 for f in path.rglob("*"):453 if f.is_file():454 try:455 total += f.stat().st_size456 except OSError:457 pass458 return total459460 usage = shutil.disk_usage(settings.models_dir)461 models_bytes, logs_bytes, db_bytes = await asyncio.gather(462 asyncio.to_thread(du, settings.models_dir), asyncio.to_thread(du, settings.logs_path),463 asyncio.to_thread(lambda: sum(du(p) for p in [settings.db_path] if p.exists()) + du(settings.data_path)))464 hf_cache = Path.home() / ".cache" / "huggingface"465 cache_bytes = await asyncio.to_thread(du, hf_cache)466 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")467 return {"total_gb": round(usage.total / GB, 1), "used_gb": round(usage.used / GB, 1), "free_gb": round(usage.free / GB, 1),468 "models_gb": round(models_bytes / GB, 2), "logs_gb": round(logs_bytes / GB, 3), "database_gb": round(db_bytes / GB, 3),469 "cache_gb": round(cache_bytes / GB, 2), "other_gb": round(max(0, usage.used - models_bytes - logs_bytes - db_bytes - cache_bytes) / GB, 1),470 "model_root": str(settings.models_dir), "models": by_model,471 "min_free_gb": float(await state.db.get_setting("min_free_disk_gb", settings.min_free_disk_gb))}472473474@router.get("/system/processes")475async def system_processes(request: Request, p: Principal = Depends(require_admin)):476 import psutil477 manager = request.app.state.manager478 procs = []479 for lm in manager.loaded.values():480 try:481 pr = psutil.Process(lm.handle.pid)482 procs.append({"model_id": lm.model["id"], "pid": lm.handle.pid, "port": lm.handle.port, "runtime": lm.handle.runtime,483 "rss_gb": round(lm.handle.memory_bytes() / GB, 2), "cpu_percent": pr.cpu_percent(interval=None),484 "status": lm.status, "threads": pr.num_threads(), "created": pr.create_time()})485 except psutil.Error:486 pass487 me = psutil.Process()488 return {"workers": procs, "server": {"pid": me.pid, "rss_gb": round(me.memory_info().rss / GB, 3),489 "cpu_percent": me.cpu_percent(interval=None), "threads": me.num_threads()}}490491492@router.get("/system/health")493async def system_health_admin(request: Request):494 return await public_health(request)495496497async def public_health(request: Request) -> dict:498 state = request.app.state499 hw = detect_hardware(state.settings.models_dir)500 tel = state.metrics.last or await state.metrics.sample()501 cur = state.manager.current_model()502 return {"status": "ok", "hardware": {"chip": hw.chip, "memory_gb": hw.memory_gb, "gpu_cores": hw.gpu_cores},503 "model": {"loaded": cur is not None, "id": cur.model["id"] if cur else None,504 "loaded_models": [lm.model["id"] for lm in state.manager.loaded.values() if lm.status == "ready"]},505 "memory": {"used_gb": tel["mem_used_gb"], "available_gb": tel["mem_available_gb"], "pressure": tel["mem_pressure_level"]},506 "uptime_seconds": round(time.time() - state.metrics.started_at), "version": request.app.version}507508509@router.get("/system/metrics")510async def system_metrics(request: Request, p: Principal = Depends(require_admin), minutes: int = Query(60, ge=5, le=60 * 24 * 7),511 hours: int = Query(24, ge=1, le=24 * 30)):512 state = request.app.state513 return {"history": await state.metrics.history(minutes), "requests": await state.metrics.request_stats(hours),514 "live": state.metrics.last}515516517@router.get("/runtime/status")518async def runtime_status(request: Request, p: Principal = Depends(require_admin)):519 return request.app.state.manager.snapshot()520521522@router.get("/runtime/current-model")523async def runtime_current(request: Request, p: Principal = Depends(require_admin)):524 cur = request.app.state.manager.current_model()525 return {"model": cur.to_dict() if cur else None}526527528# ---------------------------------------------------------------------------529# Jobs / downloads / harvest530# ---------------------------------------------------------------------------531532533@router.get("/jobs")534async def jobs(request: Request, p: Principal = Depends(require_admin), kind: str | None = None, limit: int = 100):535 kinds = {kind} if kind else None536 return {"jobs": await request.app.state.jobs.history(kinds, limit)}537538539@router.post("/jobs/{job_id}/cancel")540async def cancel_job(job_id: str, request: Request, p: Principal = Depends(require_admin)):541 ok = await request.app.state.jobs.cancel(job_id)542 return {"ok": ok}543544545class InspectBody(BaseModel):546 repository: str547 quant: str | None = None548549550class DownloadBody(BaseModel):551 repository: str552 quant: str | None = None553 force: bool = False554555556@router.post("/models/inspect")557async def inspect_repo(body: InspectBody, request: Request, p: Principal = Depends(require_admin)):558 return await request.app.state.downloader.inspect(body.repository, body.quant)559560561@router.post("/models/download")562async def download(body: DownloadBody, request: Request, p: Principal = Depends(require_admin)):563 job = await request.app.state.downloader.start_download(body.repository, body.quant, force=body.force, actor=_actor(request))564 return {"job": job.to_dict()}565566567@router.get("/downloads")568async def downloads(request: Request, p: Principal = Depends(require_admin)):569 return {"downloads": await request.app.state.jobs.history({"download"}, 100)}570571572@router.post("/downloads/{job_id}/retry")573async def retry_download(job_id: str, request: Request, p: Principal = Depends(require_admin)):574 rows = await request.app.state.db.fetchone("SELECT payload FROM jobs WHERE id=?", (job_id,))575 if not rows:576 raise APIError("Job not found.", status_code=404, code="JOB_NOT_FOUND")577 payload = json.loads(rows["payload"] or "{}")578 job = await request.app.state.downloader.start_download(payload["repository"], payload.get("quant"), force=True, actor=_actor(request))579 return {"job": job.to_dict()}580581582class HarvestBody(BaseModel):583 runtimes: list[str] | None = None584 authors: list[str] | None = None585 limit_per_author: int = 150586 min_downloads: int = 500587 max_ram_gb: float | None = None588 families: list[str] | None = None589 tasks: list[str] | None = None590 search: str | None = None591592593@router.post("/harvest/scan")594async def harvest_scan(body: HarvestBody, request: Request, p: Principal = Depends(require_admin)):595 job = await request.app.state.harvester.start_scan(body.model_dump(), actor=_actor(request))596 return {"job": job.to_dict()}597598599@router.get("/harvest/candidates")600async def harvest_candidates(request: Request, p: Principal = Depends(require_admin), include_duplicates: bool = False,601 task: str | None = None, runtime: str | None = None, family: str | None = None,602 size_class: str | None = None, q: str | None = None, limit: int = 300):603 h = request.app.state.harvester604 rows = await h.candidates(include_duplicates=include_duplicates, task=task, runtime=runtime, family=family,605 size_class=size_class, q=q, limit=limit)606 last = await request.app.state.db.scalar("SELECT MAX(scanned_at) FROM harvest_candidates")607 return {"candidates": rows, "last_scan": last, "starter": await h.suggest_starter()}608609610class SelectBody(BaseModel):611 repo_id: str612 selected: bool = True613614615@router.post("/harvest/select")616async def harvest_select(body: SelectBody, request: Request, p: Principal = Depends(require_admin)):617 await request.app.state.harvester.select(body.repo_id, body.selected)618 return {"ok": True}619620621@router.post("/harvest/dismiss")622async def harvest_dismiss(body: SelectBody, request: Request, p: Principal = Depends(require_admin)):623 await request.app.state.harvester.dismiss(body.repo_id)624 return {"ok": True}625626627@router.post("/harvest/queue")628async def harvest_queue(request: Request, p: Principal = Depends(require_admin)):629 return {"queued": await request.app.state.harvester.queue_selected(actor=_actor(request))}630631632# ---------------------------------------------------------------------------633# API keys634# ---------------------------------------------------------------------------635636637class KeyBody(BaseModel):638 name: str = Field(min_length=1, max_length=64)639 scopes: list[str] = ["inference"]640641642@router.get("/keys")643async def list_keys(request: Request, p: Principal = Depends(require_admin)):644 rows = await request.app.state.db.fetchall(645 "SELECT id, name, prefix, scopes, created_at, last_used_at, request_count, revoked_at FROM api_keys ORDER BY created_at DESC")646 for r in rows:647 r["scopes"] = r["scopes"].split(",")648 return {"keys": rows}649650651@router.post("/keys")652async def create_key(body: KeyBody, request: Request, p: Principal = Depends(require_admin)):653 scopes = [s for s in body.scopes if s in ("inference", "admin")] or ["inference"]654 raw, row = await get_auth(request).create_key(body.name, scopes)655 await request.app.state.db.audit("key.create", actor=_actor(request), target=body.name, ip=_ip(request))656 row = dict(row)657 row["scopes"] = row["scopes"].split(",")658 return {"key": raw, "record": row}659660661class KeyPatch(BaseModel):662 name: str | None = None663664665@router.patch("/keys/{key_id}")666async def rename_key(key_id: int, body: KeyPatch, request: Request, p: Principal = Depends(require_admin)):667 if body.name:668 await request.app.state.db.execute("UPDATE api_keys SET name=? WHERE id=?", (body.name, key_id))669 return {"ok": True}670671672@router.delete("/keys/{key_id}")673async def revoke_key(key_id: int, request: Request, p: Principal = Depends(require_admin)):674 await request.app.state.db.execute("UPDATE api_keys SET revoked_at=? WHERE id=? AND revoked_at IS NULL", (time.time(), key_id))675 await request.app.state.db.audit("key.revoke", actor=_actor(request), target=str(key_id), ip=_ip(request))676 return {"ok": True}677678679# ---------------------------------------------------------------------------680# Settings / logs / events681# ---------------------------------------------------------------------------682683684@router.get("/settings")685async def get_settings_(request: Request, p: Principal = Depends(require_admin)):686 state = request.app.state687 s = state.settings688 defaults = {689 "max_model_memory_gb": s.max_model_memory_gb, "absolute_max_memory_gb": s.absolute_max_memory_gb,690 "min_free_disk_gb": s.min_free_disk_gb, "model_idle_timeout_minutes": s.model_idle_timeout_minutes,691 "max_simultaneous_models": s.max_simultaneous_models, "preferred_runtime": "mlx", "default_model": None,692 "preload_model": s.preload_model, "log_prompts": s.log_prompts, "allow_downloads": s.allow_downloads,693 "allow_gguf": s.enable_gguf, "allow_mlx": s.enable_mlx, "default_context": s.default_context,694 "default_max_tokens": s.default_max_tokens, "benchmark_max_tokens": 256, "benchmark_runs": 2,695 }696 stored = await state.db.all_settings()697 merged = {**defaults, **{k: v for k, v in stored.items() if k in defaults}}698 return {"settings": merged, "defaults": defaults, "paths": {"root": str(s.root), "models": str(s.models_dir),699 "data": str(s.data_path), "logs": str(s.logs_path), "db": str(s.db_path)},700 "hf_token_set": bool(s.hf_token), "public_url": s.public_url}701702703@router.patch("/settings")704async def patch_settings(request: Request, p: Principal = Depends(require_admin)):705 state = request.app.state706 body = await request.json()707 if not isinstance(body, dict):708 raise APIError("Body must be an object.")709 changed = {}710 hw = detect_hardware(state.settings.models_dir)711 for k, v in body.items():712 if k not in SETTINGS_KEYS:713 continue714 typ = SETTINGS_KEYS[k]715 try:716 if typ is bool:717 v = bool(v)718 elif v is None or v == "":719 v = None720 else:721 v = typ(v)722 except (TypeError, ValueError):723 raise APIError(f"Invalid value for {k}.", param=k)724 if k in ("max_model_memory_gb", "absolute_max_memory_gb") and v is not None:725 if v < 1 or v > hw.memory_gb - 4:726 raise APIError(f"{k} must be between 1 and {hw.memory_gb - 4:.0f} GB on this machine.", param=k)727 if k == "max_simultaneous_models" and v is not None and (v < 1 or v > 4):728 raise APIError("max_simultaneous_models must be 1-4.", param=k)729 changed[k] = v730 await state.db.set_setting(k, v)731 if changed:732 await state.db.audit("settings.update", actor=_actor(request), detail=changed, ip=_ip(request))733 if any(k in changed for k in ("max_model_memory_gb", "absolute_max_memory_gb")):734 await state.manager.registry.reevaluate_all()735 bus.publish("settings", changed)736 return {"changed": changed}737738739@router.get("/logs/audit")740async def audit_logs(request: Request, p: Principal = Depends(require_admin), limit: int = 200):741 rows = await request.app.state.db.fetchall("SELECT * FROM audit_logs ORDER BY created_at DESC LIMIT ?", (limit,))742 return {"logs": rows}743744745@router.get("/logs/events")746async def model_events(request: Request, p: Principal = Depends(require_admin), limit: int = 200):747 rows = await request.app.state.db.fetchall("SELECT * FROM model_events ORDER BY created_at DESC LIMIT ?", (limit,))748 return {"events": rows}749750751@router.get("/logs/requests")752async def request_logs(request: Request, p: Principal = Depends(require_admin), limit: int = 200):753 rows = await request.app.state.db.fetchall(754 "SELECT id, created_at, model_id, requested_model, endpoint, api_key_id, stream, prompt_tokens, completion_tokens, ttft_ms, "755 "total_ms, tps, load_wait_ms, status, error_code FROM inference_requests ORDER BY created_at DESC LIMIT ?", (limit,))756 return {"requests": rows}757758759@router.get("/logs/worker/{model_id}")760async def worker_log(model_id: str, request: Request, p: Principal = Depends(require_admin), lines: int = 200):761 path = request.app.state.settings.logs_path / "workers" / f"worker-{model_id}.log"762 if not path.exists():763 return {"lines": []}764 txt = path.read_text(errors="replace").splitlines()765 return {"lines": txt[-lines:]}766767768@router.get("/events")769async def events(request: Request, p: Principal = Depends(require_admin)):770 q = bus.subscribe()771 state = request.app.state772773 async def gen():774 try:775 # initial snapshot776 snap = {"seq": 0, "ts": time.time(), "type": "snapshot",777 "data": {"manager": state.manager.snapshot(), "metrics": state.metrics.last,778 "jobs": state.jobs.list(limit=30)}}779 yield bus.format_sse(snap)780 while True:781 try:782 ev = await asyncio.wait_for(q.get(), timeout=15)783 yield bus.format_sse(ev)784 except asyncio.TimeoutError:785 yield ": ping\n\n"786 if await request.is_disconnected():787 break788 finally:789 bus.unsubscribe(q)790791 return StreamingResponse(gen(), media_type="text/event-stream",792 headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})793794795@router.post("/models/{model_id}/tokenize")796async def tokenize(model_id: str, request: Request, p: Principal = Depends(require_admin)):797 manager = request.app.state.manager798 lm = manager.get_ready(model_id)799 if not lm:800 raise Conflict("Model is not loaded.")801 body = await request.json()802 if lm.handle.runtime == "mlx":803 r = await manager.client.post(f"{lm.handle.base_url}/tokenize", json=body, timeout=30)804 else:805 text = body.get("text") or " ".join(m.get("content", "") for m in body.get("messages", []) if isinstance(m.get("content"), str))806 r = await manager.client.post(f"{lm.handle.base_url}/tokenize", json={"content": text}, timeout=30)807 d = r.json()808 return {"tokens": len(d.get("tokens", [])), "max_context": lm.handle.context}809 return r.json()810