bdce2764b1
Drag, paste or pick a file in the composer. Images go to vision models as multimodal content parts; PDFs and text files have their content extracted and placed in the prompt. Verified end to end against gemma4-e4b-q8 on llama-swap: given a drawing and a text file, it named the red square and blue circle and read the number out of the document. Type is decided by inspecting the bytes, never the filename or the browser's Content-Type -- a .png full of text is stored as text. Images are downscaled to 1400px and re-encoded: a phone photo is several megabytes of base64, which is slow and a large slice of the context window. PDF text is extracted once, at upload, and stored; re-extracting per request would let a reply change because a parser was upgraded. Design points worth keeping: - Images are only sent to models an administrator has marked `vision`. This is not graceful degradation -- most endpoints reject the entire request rather than ignoring an image part. A plain text turn stays a plain string for the same reason: the list form 400s on endpoints that do not implement it. - Images reach the model as base64 data URIs, not links. A local endpoint has no route back to LLeMbas, and a hosted one has no credentials for it. - Non-images are served Content-Disposition: attachment with nosniff, so an uploaded .html can never execute in this origin. Stored names are random; the uploader's name is a label and never a path. - Uploads are unbound until the message is sent, which is what lets a file be removed beforehand. claim() only takes unclaimed rows owned by the sender, so a forged id cannot pull in someone else's file. Abandoned uploads are swept at startup. - A scanned PDF says so rather than silently contributing nothing, and truncation is declared to the model in the document tag so it can admit it did not see page 400. - "Here, look at this" with no words is a legitimate turn, so a message is only empty when it carries neither text nor files. Also fixes auto-titling, which read message["content"] as a string and would have broken on the first multimodal turn. 186 tests, ruff clean. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
440 lines
16 KiB
Python
440 lines
16 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 files as files_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(""),
|
|
file_ids: list[str] = Form(default=[]),
|
|
) -> 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()
|
|
# "Here, look at this" with no words is a legitimate turn, so an empty
|
|
# message is only empty when it carries nothing at all.
|
|
if not content and not file_ids:
|
|
return Response(status_code=status.HTTP_204_NO_CONTENT)
|
|
|
|
user_message = chat_service.create_message(db, chat, ROLE_USER, content)
|
|
if file_ids:
|
|
files_service.claim(db, ids=file_ids, user_id=user.id, message_id=user_message.id)
|
|
db.refresh(user_message)
|
|
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",
|
|
},
|
|
)
|
|
|
|
|
|
def _plain_text(content: str | list) -> str:
|
|
"""The text of a message payload, whether it is a string or content parts."""
|
|
if isinstance(content, str):
|
|
return content
|
|
if isinstance(content, list):
|
|
return " ".join(
|
|
part.get("text", "")
|
|
for part in content
|
|
if isinstance(part, dict) and part.get("type") == "text"
|
|
).strip()
|
|
return ""
|
|
|
|
|
|
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)
|
|
# A multimodal turn's content is a list of parts, not a string, so
|
|
# the text has to be picked out before it can title a chat.
|
|
first_user_text = next(
|
|
(
|
|
_plain_text(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"]
|