Files
LLeMbas/src/lembas/services/chat.py
T
Jaroslav Beneš 3b1632069c Compaction: a button, and automatically when the window fills
A long conversation eventually just stops working. Compaction summarises the
earlier turns and sends the summary in their place.

The messages are kept. They stay in the transcript behind a collapsed
divider and simply stop being part of the request, which is what makes the
button safe to press and automatic compaction safe to have at all: a summary
that came out badly is a bad turn, not a lost conversation.

Stored on the Chat, not as a synthetic Message. A synthetic row needs a
role -- `system` breaks the one-system-message rule the moment build_messages
emits it beside the harness, and user/assistant makes it a turn people can
edit, regenerate from and copy, indistinguishable from a real one in all
four places a bubble is rendered. Worse, "editing rewinds, it does not
branch" would silently delete it and leave no marker that compaction had
happened at all.

The summary goes out as a user turn and an assistant turn, not one. A
leading assistant breaks templates requiring the first non-system message to
be user; a lone leading user produces user, user whenever the kept history
starts on a user turn -- which it always does, because the cutoff lands on a
finished reply.

compacted_through_id is a plain id rather than a foreign key: migrations.py
compiles only the column type, so a REFERENCES clause would exist on a fresh
database and not on an upgraded one, and a constraint half the fleet has is
worse than none. cutoff_message validates it on every read instead, and a
rewind past the boundary clears it.

Compacting again summarises only the delta, with the previous summary
supplied to be subsumed. Re-summarising the whole chat each time grows
quadratically and eventually exceeds the window it is protecting.

Automatically at the top of _run, not in post_message: that route's contract
is to return immediately and leave the slow part to a resumable connection,
and it also means build_request is called once, after compaction, with no
second assembly path. The trigger is the last reply's recorded usage plus an
estimate of the new turn -- retrospective because true prompt_tokens are only
knowable after a response, plus the delta because otherwise fifty thousand
characters pasted into the composer overflow a window that read 90% last
turn. It never fires when the context length is unknown. It does fire on
estimated counts, which is safe here precisely because nothing is lost.

_maybe_compact never raises: a failure logs and sends the uncompacted
request. A `status` event says "Summarising earlier messages…" in the
meantime, because a silent multi-second pause before the first token is what
a hang looks like.

The wording is three fragments under Admin - Prompts. Clearing task.compact
turns compaction off entirely.

Also adds compaction.moment(): SQLite does not store the offset, so a row
loaded from disk is naive while one in the session's identity map keeps its
tzinfo, and comparing the two raises. Every comparison here is between
exactly those.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-01 01:02:02 +02:00

514 lines
18 KiB
Python

"""Chat orchestration: building requests, streaming replies, naming chats."""
from __future__ import annotations
import logging
from datetime import UTC, datetime, timedelta
from typing import Any
from sqlalchemy import func, select
from sqlalchemy.orm import Session as DBSession
from lembas.db.models import (
ROLE_ASSISTANT,
ROLE_SYSTEM,
ROLE_USER,
Chat,
Connection,
Message,
Model,
)
from lembas.services import files as files_service
from lembas.services.llm.openai_client import Endpoint, LLMError, complete
log = logging.getLogger(__name__)
# Sampling keys forwarded upstream. Anything else a user puts in params_json is
# ignored rather than passed through, so a typo cannot produce a 400 from the
# provider that looks like a LLeMbas bug.
FORWARDED_PARAMS = frozenset(
{"temperature", "top_p", "max_tokens", "presence_penalty", "frequency_penalty",
"seed", "stop"}
)
MAX_TITLE_LENGTH = 60
# How long a temporary chat survives after the last thing said in it.
TEMPORARY_LIFETIME = timedelta(hours=24)
def resolve_endpoint(db: DBSession, chat: Chat) -> tuple[Endpoint, str]:
"""Find the connection and model a chat should use.
Chats store the model id as text rather than a foreign key so history
survives an admin deleting a connection, which means the mapping back to a
live connection has to be resolved at send time and can legitimately fail.
"""
if not chat.model_id:
raise LLMError("This chat has no model selected.")
connection: Connection | None = None
if chat.connection_id:
connection = db.get(Connection, chat.connection_id)
if connection is None or not connection.enabled:
# The original connection is gone or disabled. Any enabled connection
# still offering this model id will do.
model = db.scalar(
select(Model)
.join(Connection)
.where(
Model.model_id == chat.model_id,
Model.enabled.is_(True),
Connection.enabled.is_(True),
)
.order_by(Connection.position)
)
if model is None:
raise LLMError(
f"No enabled connection currently offers the model "
f"'{chat.model_id}'. Pick another model for this chat."
)
connection = model.connection
chat.connection_id = connection.id
db.commit()
return Endpoint.from_connection(connection), chat.model_id
def document_context(message: Message) -> str:
"""Extracted text from a message's non-image attachments.
Wrapped in named tags so the model can tell one document from another, and
tell all of them from what the user actually typed. Truncation is stated
inline rather than silently, so a model asked about page 400 of a 300-page
extract can say it did not see it.
"""
blocks: list[str] = []
for attachment in message.documents:
if not attachment.extracted_text.strip():
continue
note = " (truncated)" if attachment.truncated else ""
blocks.append(
f'<document name="{attachment.filename}"{note}>\n'
f"{attachment.extracted_text.strip()}\n"
f"</document>"
)
return "\n\n".join(blocks)
def message_payload(message: Message, *, vision: bool) -> dict[str, Any]:
"""One history entry in the shape the endpoint expects.
Plain text stays a plain string: sending the multimodal list form to an
endpoint that does not implement it is a reliable way to get a 400, and
most local runners do not.
"""
text = message.content.strip()
documents = document_context(message)
if documents:
# Documents lead so the question that follows has its material already
# in view, which is how these models are trained to read a prompt.
text = f"{documents}\n\n{text}" if text else documents
images = message.images if vision else []
if not images:
return {"role": message.role, "content": text}
parts: list[dict[str, Any]] = []
if text:
parts.append({"type": "text", "text": text})
for attachment in images:
uri = files_service.data_uri(attachment)
if uri is None:
# The row survived but the file did not. Better to say so than to
# send a turn that silently lost its picture.
log.warning("attachment %s has no file on disk", attachment.id)
continue
parts.append({"type": "image_url", "image_url": {"url": uri}})
if not parts:
return {"role": message.role, "content": text}
return {"role": message.role, "content": parts}
def effective_system_prompt(db: DBSession, chat: Chat) -> str:
"""The system prompt a chat actually runs with.
Three layers, most specific wins outright:
chat > model > instance
Precedence rather than concatenation. Stacking them reads well in a
settings screen and badly in practice: the moment two layers disagree the
model gets contradictory instructions and nobody can tell which one is
losing. With precedence, "why is it behaving like this" has one answer.
"""
from lembas.services import settings_store
if chat.system_prompt.strip():
return chat.system_prompt.strip()
model = db.scalar(
select(Model).where(Model.model_id == chat.model_id).order_by(Model.position)
)
if model is not None and (model.system_prompt or "").strip():
return model.system_prompt.strip()
return (settings_store.get(db, "system_prompt") or "").strip()
def build_messages(
db: DBSession,
chat: Chat,
*,
upto: Message | None = None,
vision: bool = False,
system_prompt: str | None = None,
) -> list[dict]:
"""Assemble the message list to send upstream.
`upto` excludes the placeholder assistant row being generated into, and
everything after it. `system_prompt` overrides what would otherwise be
resolved, which is how the harness gets in front of the authored prompt
without this function knowing anything about tools.
"""
from lembas.services import compaction as compaction_service
from lembas.services import prompts as prompts_service
payload: list[dict[str, Any]] = []
system = effective_system_prompt(db, chat) if system_prompt is None else system_prompt
if system:
payload.append({"role": ROLE_SYSTEM, "content": system})
# Compacted turns are replaced by a summary carried in two turns rather than
# one. A leading `assistant` breaks templates that require the first
# non-system message to be `user`; a lone leading `user` produces user, user
# whenever the kept history starts on a user turn -- which it always does,
# because the cutoff lands on a finished reply. The pair alternates
# correctly in both directions and keeps exactly one system message.
cutoff = compaction_service.cutoff_message(db, chat)
if cutoff is not None:
lead = prompts_service.resolve(db, "task.compact_lead").strip()
ack = prompts_service.resolve(db, "task.compact_ack").strip()
summary = chat.compact_summary.strip()
payload.append(
{"role": ROLE_USER, "content": f"{lead}\n\n{summary}" if lead else summary}
)
if ack:
payload.append({"role": ROLE_ASSISTANT, "content": ack})
history = db.scalars(
select(Message).where(Message.chat_id == chat.id).order_by(Message.created_at)
).all()
for message in history:
if upto is not None and message.id == upto.id:
break
if cutoff is not None and compaction_service.moment(
message
) <= compaction_service.moment(cutoff):
continue
# Skip turns that failed or produced nothing -- but a message carrying
# only an attachment has no text and must still be sent.
if message.error:
continue
if not message.content.strip() and not message.attachments:
continue
payload.append(message_payload(message, vision=vision))
return payload
def model_for(db: DBSession, chat: Chat) -> Model | None:
"""The Model row a chat is using, or None if it has gone.
Looked up by id rather than held as a foreign key, for the same reason
resolve_endpoint does: chats store the model as text so history survives an
administrator deleting a connection.
"""
return db.scalar(
select(Model).where(Model.model_id == chat.model_id).order_by(Model.position)
)
def model_supports(db: DBSession, chat: Chat, capability: str) -> bool:
"""Whether the chat's current model is marked as having a capability."""
model = model_for(db, chat)
return bool(model and (model.capabilities_json or {}).get(capability))
def build_request(
db: DBSession,
chat: Chat,
*,
upto: Message | None = None,
tools: list[dict[str, Any]] | None = None,
user=None,
) -> dict[str, Any]:
"""The whole request body, tools and harness included.
Composed here rather than in the generation loop so that "what gets sent"
has one answer, and so the harness cannot be forgotten by a future caller
that offers tools.
"""
from lembas.services import harness as harness_service
from lembas.services import prompts as prompts_service
params = {
key: value
for key, value in (chat.params_json or {}).items()
if key in FORWARDED_PARAMS and value not in (None, "")
}
# Images are only sent to a model an administrator has marked as having
# vision. Sending them to one that has not is not a graceful degradation:
# most endpoints reject the whole request.
vision = model_supports(db, chat, "vision")
if user is None:
from lembas.db.models import User
user = db.get(User, chat.user_id)
# The harness describes the tools; the authored prompt describes the
# behaviour. See services/harness.py for why these are joined rather than
# being two competing layers.
system = harness_service.join(
harness_service.compose(db, user, tools, chat),
effective_system_prompt(db, chat),
lead=prompts_service.render(db, "seam.authored_lead", {}),
)
body: dict[str, Any] = {
"model": chat.model_id,
"messages": build_messages(
db, chat, upto=upto, vision=vision, system_prompt=system
),
**params,
}
if tools:
body["tools"] = tools
return body
def default_model(db: DBSession, user=None) -> tuple[str, str] | None:
"""The model a new chat should start with, as (model_id, connection_id).
Preference order: the user's own choice, then the instance default, then
whatever is first in the admin's ordering. Each is checked against what the
user may actually reach, so a default they have lost access to falls
through rather than producing a chat they cannot use.
"""
from lembas.security import permissions
from lembas.services import settings_store
reachable = permissions.models_visible_to(db, user)
if not reachable:
return None
by_id = {model.model_id: model for model in reachable}
preferred = (user.settings_json or {}).get("default_model") if user is not None else None
if preferred and preferred in by_id:
return preferred, by_id[preferred].connection_id
instance_default = settings_store.get(db, "default_model")
if instance_default and instance_default in by_id:
return instance_default, by_id[instance_default].connection_id
# First in the administrator's ordering. Pinning is a sidebar shortcut, not
# a reordering, so it deliberately does not influence this.
chosen = sorted(reachable, key=lambda m: (m.position, m.model_id))[0]
return chosen.model_id, chosen.connection_id
def available_models(db: DBSession, user=None) -> list[Model]:
"""Models this user may start a chat with, in the administrator's order.
Pinning does NOT hoist a model up this list: pinned models get their own
shortcuts in the sidebar, and a picker whose order silently differs from
the one configured in the admin screen is just confusing.
"""
from lembas.security import permissions
reachable = permissions.models_visible_to(db, user)
return sorted(reachable, key=lambda m: (m.position, m.model_id))
def fallback_title(text: str) -> str:
"""Derive a chat title from the opening message, without calling a model."""
cleaned = " ".join(text.split())
if not cleaned:
return "New chat"
if len(cleaned) <= MAX_TITLE_LENGTH:
return cleaned
# Prefer a word boundary, but only if it does not cut the title in half.
clipped = cleaned[:MAX_TITLE_LENGTH]
space = clipped.rfind(" ")
if space > MAX_TITLE_LENGTH * 0.6:
clipped = clipped[:space]
return clipped.rstrip(" ,.;:-") + "…"
async def generate_title(
endpoint: Endpoint, model_id: str, question: str, answer: str, *, template: str
) -> str:
"""Ask the model for a short chat title.
Best-effort by design: any failure falls back to trimming the first
message. Naming a chat is never worth surfacing an error for.
`template` is passed in rather than read here because this runs after the
generation's session has closed -- see `generation._run`. An empty one means
an administrator cleared the fragment, which is how auto-titling is turned
off: no request is made at all.
"""
from lembas.services import prompts as prompts_service
if not template.strip():
return fallback_title(question)
prompt = prompts_service.substitute(
template, {"question": question[:500], "answer": answer[:500]}
)
try:
raw = await complete(
endpoint,
{
"model": model_id,
"messages": [{"role": ROLE_USER, "content": prompt}],
"max_tokens": 24,
"temperature": 0.2,
},
)
except LLMError as exc:
log.debug("auto-title failed, using fallback: %s", exc)
return fallback_title(question)
title = " ".join(raw.split()).strip().strip('"“”\'')
# Small models sometimes ignore the instruction and answer the question
# instead; an over-long reply is a better signal of that than anything else.
if not title or len(title) > MAX_TITLE_LENGTH * 1.5:
return fallback_title(question)
return title[:MAX_TITLE_LENGTH]
def create_message(
db: DBSession,
chat: Chat,
role: str,
content: str = "",
*,
complete_: bool = True,
model_id: str = "",
) -> Message:
message = Message(
chat_id=chat.id,
role=role,
content=content,
complete=complete_,
model_id=model_id,
)
db.add(message)
db.commit()
return message
async def summarise_for_compaction(
endpoint: Endpoint,
model_id: str,
*,
transcript: str,
previous_summary: str,
template: str,
) -> str:
"""Ask the model to summarise the earlier turns.
`template` is passed in for the same reason `generate_title`'s is: this runs
after the generation's session has closed, and opening another one there is
how you get a session that outlives its scope. An empty template means an
administrator cleared the fragment, and nothing is asked of anyone.
"""
from lembas.services import prompts as prompts_service
if not template.strip() or not transcript.strip():
return ""
prompt = prompts_service.substitute(
template, {"transcript": transcript, "previous_summary": previous_summary}
)
raw = await complete(
endpoint,
{
"model": model_id,
"messages": [{"role": ROLE_USER, "content": prompt}],
"max_tokens": 1200,
# Low, but not zero: this is recall, not invention.
"temperature": 0.3,
},
)
return raw.strip()
def sweep_temporary(db: DBSession, older_than: timedelta = TEMPORARY_LIFETIME) -> int:
"""Delete temporary chats nobody has touched for a day.
Age is measured from the newest message rather than from the chat row's own
timestamps. `created_at` would destroy a conversation still in use at hour
23, and `updated_at` does not move when a message is inserted -- `onupdate`
fires on an UPDATE of the chat, and adding a message is not one.
Startup only, like files.sweep_orphans beside it. A server that runs for a
month sweeps once; that is the trade the existing sweep already makes, and a
scheduler is a whole new concern for a single-worker application.
"""
cutoff = datetime.now(UTC) - older_than
newest = (
select(Message.chat_id, func.max(Message.created_at).label("last"))
.group_by(Message.chat_id)
.subquery()
)
stale = list(
db.scalars(
select(Chat)
.outerjoin(newest, newest.c.chat_id == Chat.id)
.where(
Chat.temporary.is_(True),
func.coalesce(newest.c.last, Chat.created_at) < cutoff,
)
)
)
if not stale:
return 0
files_service.remove_files_for_chats(db, [chat.id for chat in stale])
for chat in stale:
db.delete(chat)
db.commit()
log.info("swept %d temporary chat(s)", len(stale))
return len(stale)
def user_chats(db: DBSession, user_id: str, *, folder_id: str | None = None) -> list[Chat]:
query = select(Chat).where(
Chat.user_id == user_id, Chat.archived.is_(False), Chat.temporary.is_(False)
)
if folder_id is not None:
query = query.where(Chat.folder_id == folder_id)
return list(db.scalars(query.order_by(Chat.pinned.desc(), Chat.updated_at.desc())))
__all__ = [
"ROLE_ASSISTANT",
"ROLE_USER",
"available_models",
"build_request",
"create_message",
"default_model",
"fallback_title",
"generate_title",
"resolve_endpoint",
"user_chats",
]