SPB Git forge

spb/llm-api

Public
0commits 0branches 0releases
0 Bsize
maindefault branch
—last push
33.8 KB · 810 lines python
Raw Blame History
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