Files
LLeMbas/src/lembas/api/chats.py
T
Jaroslav Beneš d6c87ac811 Users, groups, permissions, model settings and reasoning display
Four features, plus the schema machinery they needed.

**Schema sync.** The first live instance had data in it, and create_all
only creates missing *tables* -- a new column silently never appeared.
db/migrations.py now diffs the declared models against the database and
ALTER TABLE ... ADD COLUMN for what is missing, deriving a backfill
default from the column type (SQLite refuses a NOT NULL column without
one, and a Python-side `default=dict` cannot be expressed in DDL).
Verified against a copy of the live database: eight changes applied, all
rows preserved, second run a no-op. Renames, drops and retypes are still
manual and say so.

**Permissions.** A flat set of named booleans: an instance baseline
widened by each group the user belongs to. A group grants and never
denies -- with denies, "why can this user not do X" cannot be answered
without simulating every group. Admins bypass entirely, because an admin
can grant it back to themselves in two clicks and pretending otherwise
is theatre. Model *access* is separate: public, or granted to groups.
The picker is not the boundary -- switching a chat to a model you cannot
reach is a 403.

**Model settings.** Ordering, pinned-first, an instance default and a
per-user default, display names, descriptions, capability flags, and
uploaded images. Images are stored and served locally rather than by
URL: a remote URL makes every page render a request to a third party.
Uploads are validated by magic number, not the declared content type,
and stored under a random name. Models with no image get a generated
initial whose hue is derived from the model id, so it is stable.

**Reasoning display.** Streams into its own collapsible block above the
answer, labelled "Thought for 14 seconds", collapsed once finished, and
never replayed as context on the next turn. Two sources: the
reasoning_content delta field, and <think> tags inline in content -- the
latter needs a streaming splitter because the tags arrive split across
chunks. Models emitting no reasoning show nothing, via a :has() rule
rather than JavaScript. Verified against qwen35-9b on llama-swap: 694
reasoning events, 52 answer tokens, cleanly separated.

Two bugs found and fixed while testing:

- A bare `Mapped[list]` relationship is treated by SQLAlchemy as a scalar
  and returns None instead of []. It needs the element type.
- FastAPI substitutes the default for an empty form value, so with
  `x: str | None = Form(None)` a submitted `x=` is indistinguishable from
  an absent field. That silently broke clearing a system prompt or a
  temperature. update_chat now reads the raw form and checks key presence.

143 tests, ruff clean.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-21 11:49:32 +02:00

414 lines
15 KiB
Python

"""Chat creation, messaging and the streaming reply endpoint."""
from __future__ import annotations
import asyncio
import logging
import time
from collections.abc import AsyncIterator
from fastapi import APIRouter, Depends, Form, HTTPException, Request, status
from fastapi.responses import HTMLResponse, Response, StreamingResponse
from sqlalchemy.orm import Session as DBSession
from lembas.api.deps import Db, RequiredUser, require_permission
from lembas.db.models import ROLE_ASSISTANT, ROLE_USER, Chat, Message, User
from lembas.db.session import session_scope
from lembas.security import permissions
from lembas.services import chat as chat_service
from lembas.services import sse
from lembas.services.llm.openai_client import (
LLMError,
delta_reasoning,
delta_text,
stream_chat,
)
from lembas.services.markdown import escape_text, render_markdown
from lembas.services.reasoning import REASONING, ReasoningSplitter
from lembas.web.templating import render, templates
log = logging.getLogger(__name__)
router = APIRouter(prefix="/api/chats", tags=["chats"])
def _owned_chat(db: DBSession, chat_id: str, user_id: str) -> Chat:
chat = db.get(Chat, chat_id)
# 404 rather than 403 for someone else's chat: whether a given id exists is
# not information this endpoint should hand out.
if chat is None or chat.user_id != user_id:
raise HTTPException(status.HTTP_404_NOT_FOUND, "That chat no longer exists.")
return chat
@router.post("", dependencies=[Depends(require_permission("chat.create"))])
async def create_chat(db: Db, user: RequiredUser, folder_id: str = Form("")) -> Response:
chosen = chat_service.default_model(db, user)
chat = Chat(
user_id=user.id,
folder_id=folder_id or None,
model_id=chosen[0] if chosen else "",
connection_id=chosen[1] if chosen else None,
)
db.add(chat)
db.commit()
# HX-Redirect rather than a swap: a new chat is a new URL, and the address
# bar has to follow so the chat can be reloaded or bookmarked.
response = Response(status_code=status.HTTP_204_NO_CONTENT)
response.headers["HX-Redirect"] = f"/chat/{chat.id}"
return response
@router.post("/{chat_id}/messages")
async def post_message(
request: Request,
db: Db,
user: RequiredUser,
chat_id: str,
content: str = Form(...),
) -> Response:
"""Persist the user's turn and hand back the pair of bubbles.
The assistant bubble comes back empty, carrying the sse-connect attribute
that opens the stream below. Splitting it this way means the POST returns
immediately and the slow part is a separate, resumable connection.
"""
chat = _owned_chat(db, chat_id, user.id)
content = content.strip()
if not content:
return Response(status_code=status.HTTP_204_NO_CONTENT)
user_message = chat_service.create_message(db, chat, ROLE_USER, content)
assistant_message = chat_service.create_message(
db, chat, ROLE_ASSISTANT, "", complete_=False, model_id=chat.model_id
)
# `user` is required by the shared message template, which renders both
# roles; without it the user bubble's initial blows up.
return templates.TemplateResponse(
request,
"chat/_turn.html",
{
"request": request,
"user_message": user_message,
"assistant_message": assistant_message,
"chat": chat,
"user": user,
},
)
@router.get("/{chat_id}/messages/{message_id}/stream")
async def stream_message(
db: Db,
user: RequiredUser,
chat_id: str,
message_id: str,
) -> Response:
"""Stream the assistant's reply as server-sent events.
Emits `token` events carrying escaped text, then a single `done` event
carrying the finished bubble rendered from Markdown, then `close`.
"""
chat = _owned_chat(db, chat_id, user.id)
message = db.get(Message, message_id)
if message is None or message.chat_id != chat.id:
raise HTTPException(status.HTTP_404_NOT_FOUND, "That message no longer exists.")
return StreamingResponse(
_generate(chat.id, message.id),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache, no-transform",
"Connection": "keep-alive",
# nginx buffers proxied responses by default, which turns a stream
# into one delivery at the end. This is the documented opt-out.
"X-Accel-Buffering": "no",
},
)
async def _generate(chat_id: str, message_id: str) -> AsyncIterator[str]:
"""Drive one completion and frame it as SSE.
Opens its own database session rather than using the request's: streaming
outlives the request handler, and the dependency-scoped session may already
be closed by the time the first token arrives.
"""
accumulated: list[str] = []
thinking: list[str] = []
error: str | None = None
reasoning_ms = 0
with session_scope() as db:
chat = db.get(Chat, chat_id)
message = db.get(Message, message_id)
if chat is None or message is None:
yield sse.event("close", "")
return
first_user_text = ""
# Handles models that emit <think> tags inline in content rather than
# using the reasoning_content field.
splitter = ReasoningSplitter()
started = time.monotonic()
reasoning_started: float | None = None
try:
endpoint, model_id = chat_service.resolve_endpoint(db, chat)
payload = chat_service.build_request(db, chat, upto=message)
first_user_text = next(
(m["content"] for m in reversed(payload["messages"]) if m["role"] == ROLE_USER),
"",
)
async for chunk in stream_chat(endpoint, payload):
# A dedicated reasoning field is unambiguous; take it as-is.
thought = delta_reasoning(chunk)
if thought:
if reasoning_started is None:
reasoning_started = time.monotonic()
thinking.append(thought)
yield sse.event("reasoning", escape_text(thought))
await asyncio.sleep(0)
text = delta_text(chunk)
if not text:
continue
for kind, piece in splitter.feed(text):
if kind == REASONING:
if reasoning_started is None:
reasoning_started = time.monotonic()
thinking.append(piece)
yield sse.event("reasoning", escape_text(piece))
else:
# First answer token ends the thinking phase.
if reasoning_started is not None and not reasoning_ms:
reasoning_ms = int((time.monotonic() - reasoning_started) * 1000)
accumulated.append(piece)
yield sse.event("token", escape_text(piece))
# Hand control back so the event is flushed rather than
# batched behind a fast generator.
await asyncio.sleep(0)
for kind, piece in splitter.flush():
if kind == REASONING:
thinking.append(piece)
yield sse.event("reasoning", escape_text(piece))
else:
accumulated.append(piece)
yield sse.event("token", escape_text(piece))
except LLMError as exc:
error = exc.message
log.info("generation failed for chat %s: %s", chat_id, exc.message)
except asyncio.CancelledError:
# The reader navigated away or closed the tab. Keep whatever was
# produced so the partial reply is still there on reload.
message.content = "".join(accumulated)
message.reasoning = "".join(thinking)
message.complete = True
db.commit()
raise
except Exception as exc: # noqa: BLE001 - must not kill the stream silently
error = "Something went wrong while generating this reply."
log.exception("unexpected generation failure for chat %s: %s", chat_id, exc)
if reasoning_started is not None and not reasoning_ms:
# Reasoning ran to the end without an answer following it.
reasoning_ms = int((time.monotonic() - reasoning_started) * 1000)
message.content = "".join(accumulated)
message.reasoning = "".join(thinking)
message.reasoning_ms = reasoning_ms
message.error = error or ""
message.complete = True
log.debug(
"chat %s: %d chars answer, %d chars reasoning, %.1fs total",
chat_id,
len(message.content),
len(message.reasoning),
time.monotonic() - started,
)
if not chat.title_generated and (accumulated or error):
chat.title = (
await chat_service.generate_title(
endpoint, model_id, first_user_text, message.content
)
if not error and first_user_text
else chat_service.fallback_title(first_user_text)
)
chat.title_generated = True
db.commit()
final_html = templates.get_template("chat/_message.html").render(
{
"message": message,
"body_html": render_markdown(message.content),
"chat": chat,
# Passed even though an assistant bubble never reads it: the
# template shares both roles, and a missing `user` would only
# blow up on whichever branch is not being exercised here.
"user": db.get(User, chat.user_id),
}
)
title_html = templates.get_template("chat/_title_oob.html").render(
{"chat": chat}
)
yield sse.event("done", final_html + title_html)
yield sse.event("close", "")
@router.patch("/{chat_id}")
async def update_chat(request: Request, db: Db, user: RequiredUser, chat_id: str) -> Response:
"""Partially update a chat.
The raw form is read rather than declaring Form() parameters because
FastAPI substitutes the default for an empty form value, which makes
"field absent" and "field submitted empty" indistinguishable. That
difference is exactly what this endpoint needs: an empty system prompt or
temperature means *clear it*, not *leave it alone*.
"""
chat = _owned_chat(db, chat_id, user.id)
allowed = permissions.resolve(db, user)
form = await request.form()
if "title" in form:
cleaned = str(form["title"]).strip()[:300]
if cleaned:
chat.title = cleaned
# An explicit rename must not be overwritten by auto-titling later.
chat.title_generated = True
if "folder_id" in form:
chat.folder_id = str(form["folder_id"]) or None
model_id = str(form.get("model_id", "")).strip()
if model_id:
if not allowed.get("chat.model_select"):
raise HTTPException(
status.HTTP_403_FORBIDDEN, "You may not change the model for a chat."
)
# Checked against what this user can reach, not merely what exists --
# otherwise the picker is advisory and a crafted request bypasses it.
match = next(
(m for m in chat_service.available_models(db, user) if m.model_id == model_id),
None,
)
if match is None:
raise HTTPException(status.HTTP_403_FORBIDDEN, "That model is not available to you.")
chat.model_id = model_id
chat.connection_id = match.connection_id
if "system_prompt" in form:
if not allowed.get("chat.system_prompt"):
raise HTTPException(
status.HTTP_403_FORBIDDEN, "You may not set a system prompt."
)
chat.system_prompt = str(form["system_prompt"]).strip()[:8000]
submitted_params = {name: form[name] for name in _PARAM_RANGES if name in form}
if submitted_params:
if not allowed.get("chat.params"):
raise HTTPException(
status.HTTP_403_FORBIDDEN, "You may not change sampling parameters."
)
chat.params_json = {
**(chat.params_json or {}),
**_clean_params(**{k: str(v) for k, v in submitted_params.items()}),
}
db.commit()
return Response(status_code=status.HTTP_204_NO_CONTENT)
# Bounds are the ones every provider agrees on. Out-of-range values are
# dropped rather than clamped: silently changing what someone typed is worse
# than ignoring it, and the form shows what actually stuck on reload.
_PARAM_RANGES: dict[str, tuple[type, float, float]] = {
"temperature": (float, 0.0, 2.0),
"top_p": (float, 0.0, 1.0),
"max_tokens": (int, 1, 1_000_000),
}
def _clean_params(**submitted: str | None) -> dict[str, float | int | None]:
"""Parse sampling parameters, dropping anything unusable.
An empty string means "unset this and let the provider default apply", so
it maps to None rather than being ignored.
"""
cleaned: dict[str, float | int | None] = {}
for name, raw in submitted.items():
if raw is None:
continue
if not raw.strip():
cleaned[name] = None
continue
caster, low, high = _PARAM_RANGES[name]
try:
value = caster(raw)
except (TypeError, ValueError):
continue
if low <= value <= high:
cleaned[name] = value
return cleaned
@router.delete("/{chat_id}", dependencies=[Depends(require_permission("chat.delete"))])
async def delete_chat(db: Db, user: RequiredUser, chat_id: str) -> Response:
chat = _owned_chat(db, chat_id, user.id)
db.delete(chat)
db.commit()
response = Response(status_code=status.HTTP_204_NO_CONTENT)
response.headers["HX-Redirect"] = "/chat"
return response
@router.get("/{chat_id}/messages/{message_id}/raw")
async def raw_message(db: Db, user: RequiredUser, chat_id: str, message_id: str) -> HTMLResponse:
"""The unrendered Markdown of a message, for the copy button."""
_owned_chat(db, chat_id, user.id)
message = db.get(Message, message_id)
if message is None or message.chat_id != chat_id:
raise HTTPException(status.HTTP_404_NOT_FOUND, "That message no longer exists.")
return HTMLResponse(escape_text(message.content))
@router.post("/{chat_id}/messages/{message_id}/regenerate")
async def regenerate(
request: Request,
db: Db,
user: RequiredUser,
chat_id: str,
message_id: str,
) -> Response:
"""Discard an assistant reply and produce a fresh one in its place."""
chat = _owned_chat(db, chat_id, user.id)
message = db.get(Message, message_id)
if message is None or message.chat_id != chat.id or message.role != ROLE_ASSISTANT:
raise HTTPException(status.HTTP_404_NOT_FOUND, "That reply no longer exists.")
message.content = ""
message.error = ""
message.complete = False
message.model_id = chat.model_id
db.commit()
return templates.TemplateResponse(
request,
"chat/_message.html",
{"request": request, "message": message, "chat": chat, "body_html": "", "user": user},
)
__all__ = ["render", "router"]