Data groups: a provider's models read only their own group's data
Every connection is in a data group. Its models are handed, and can find, only that group's memories, notes, skills, knowledge, reports and personality -- by search and by id. A chat stays in the group it was started in: switching its model, the endpoint fallback, the crowd, friends, bases and the @ menu all stay inside it, and a chat whose model has moved is refused rather than sent. A group may name its own embedder and image reviewer. data.manage lets a person make personal groups, remap connections for themselves and move their own records. Also: a search no longer mixes two embedders of the same width. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -32,7 +32,7 @@ from typing import Any
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.orm import Session as DBSession
|
||||
|
||||
from lembas.db.models import AUTHOR_MODEL, KIND_TASK, SOURCE_CHAT, Chat, User
|
||||
from lembas.db.models import AUTHOR_MODEL, DEFAULT_GROUP, KIND_TASK, SOURCE_CHAT, Chat, User
|
||||
from lembas.db.session import session_scope
|
||||
from lembas.services import personas as personas_service
|
||||
from lembas.services import prompts as prompts_service
|
||||
@@ -289,6 +289,12 @@ class ToolContext:
|
||||
image_checkpoint: str = ""
|
||||
model_id: str = ""
|
||||
connection_id: str = ""
|
||||
# Which data group this call reads and writes, resolved from the answering
|
||||
# model's connection. Every library runner passes it to its store and every
|
||||
# write stamps it, so a model reaches exactly one group's records -- by
|
||||
# search *and* by id, since a model that learned an id from somewhere else
|
||||
# must not be able to fetch the record past the filter.
|
||||
data_group: str = DEFAULT_GROUP
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -465,12 +471,18 @@ async def _run_knowledge_search(context: ToolContext, args: dict[str, Any]) -> T
|
||||
# session held across one is the trade `_maybe_compact` already refuses.
|
||||
# None for every "no" -- no model configured, endpoint down -- and the
|
||||
# search is then exactly the keyword one it has always been.
|
||||
vector = await _query_vector(query)
|
||||
vector = await _query_vector(query, context.data_group)
|
||||
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
found = documents_service.search(
|
||||
db, user, query, limit=6, base_ids=context.base_ids, vector=vector
|
||||
db,
|
||||
user,
|
||||
query,
|
||||
limit=6,
|
||||
base_ids=context.base_ids,
|
||||
vector=vector,
|
||||
group=context.data_group,
|
||||
)
|
||||
event = {
|
||||
"name": "knowledge_search",
|
||||
@@ -502,7 +514,7 @@ async def _run_knowledge_get(context: ToolContext, args: dict[str, Any]) -> Tool
|
||||
document_id = str(args.get("id") or "").strip()
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
document = documents_service.get(db, document_id, user)
|
||||
document = documents_service.get(db, document_id, user, context.data_group)
|
||||
if document is None:
|
||||
return ToolOutcome(
|
||||
"There is no such document, or it is not available to you.",
|
||||
@@ -526,7 +538,7 @@ async def _run_knowledge_get(context: ToolContext, args: dict[str, Any]) -> Tool
|
||||
return ToolOutcome(f"{document.title}\n\n{body}", event)
|
||||
|
||||
|
||||
async def _query_vector(query: str) -> list[float] | None:
|
||||
async def _query_vector(query: str, group: str | None = None) -> list[float] | None:
|
||||
"""The query as a vector, for the stores that can use one.
|
||||
|
||||
Its own session, opened and closed before the caller opens theirs: this is
|
||||
@@ -538,20 +550,22 @@ async def _query_vector(query: str) -> list[float] | None:
|
||||
from lembas.services.library import retrieval
|
||||
|
||||
with session_scope() as db:
|
||||
worker = retrieval.worker_for(db)
|
||||
worker = retrieval.worker_for(db, group)
|
||||
return await retrieval.embed_with(worker, query)
|
||||
|
||||
|
||||
# --- Notes -------------------------------------------------------------------
|
||||
async def _run_notes_search(context: ToolContext, args: dict[str, Any]) -> ToolOutcome:
|
||||
query = str(args.get("query") or "").strip()
|
||||
vector = await _query_vector(query)
|
||||
vector = await _query_vector(query, context.data_group)
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
found = (
|
||||
notes_service.search(db, user, query, limit=8, vector=vector)
|
||||
notes_service.search(
|
||||
db, user, query, limit=8, vector=vector, group=context.data_group
|
||||
)
|
||||
if query
|
||||
else notes_service.recent(db, user, limit=8)
|
||||
else notes_service.recent(db, user, limit=8, group=context.data_group)
|
||||
)
|
||||
event = {
|
||||
"name": "notes_search",
|
||||
@@ -571,7 +585,7 @@ async def _run_notes_search(context: ToolContext, args: dict[str, Any]) -> ToolO
|
||||
async def _run_notes_get(context: ToolContext, args: dict[str, Any]) -> ToolOutcome:
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
note = notes_service.get(db, str(args.get("id") or ""), user)
|
||||
note = notes_service.get(db, str(args.get("id") or ""), user, context.data_group)
|
||||
if note is None:
|
||||
return ToolOutcome(
|
||||
"There is no such note, or it is not available to you.",
|
||||
@@ -599,7 +613,12 @@ async def _run_notes_create(context: ToolContext, args: dict[str, Any]) -> ToolO
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
note = notes_service.create(
|
||||
db, owner=user, title=title, body=body, author=AUTHOR_MODEL
|
||||
db,
|
||||
owner=user,
|
||||
title=title,
|
||||
body=body,
|
||||
author=AUTHOR_MODEL,
|
||||
group=context.data_group,
|
||||
)
|
||||
return ToolOutcome(
|
||||
f"Saved note {note.id} — {note.title!r}.",
|
||||
@@ -615,7 +634,7 @@ async def _run_notes_create(context: ToolContext, args: dict[str, Any]) -> ToolO
|
||||
async def _run_notes_edit(context: ToolContext, args: dict[str, Any]) -> ToolOutcome:
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
note = notes_service.get(db, str(args.get("id") or ""), user)
|
||||
note = notes_service.get(db, str(args.get("id") or ""), user, context.data_group)
|
||||
if note is None or note.owner_id != context.owner_id:
|
||||
return ToolOutcome(
|
||||
"There is no such note, or it belongs to someone else. A note "
|
||||
@@ -642,7 +661,7 @@ async def _run_notes_edit(context: ToolContext, args: dict[str, Any]) -> ToolOut
|
||||
async def _run_notes_delete(context: ToolContext, args: dict[str, Any]) -> ToolOutcome:
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
note = notes_service.get(db, str(args.get("id") or ""), user)
|
||||
note = notes_service.get(db, str(args.get("id") or ""), user, context.data_group)
|
||||
if note is None or note.owner_id != context.owner_id:
|
||||
return ToolOutcome(
|
||||
"There is no such note, or it belongs to someone else.",
|
||||
@@ -735,7 +754,7 @@ async def _run_persona_write(context: ToolContext, args: dict[str, Any]) -> Tool
|
||||
return _persona_error("persona_write", "There is nobody here to be this with.")
|
||||
row = personas_service.write(
|
||||
db,
|
||||
model_key=context.model_id,
|
||||
model_key=personas_service.key_for(context.model_id, context.data_group),
|
||||
owner=user,
|
||||
content=content,
|
||||
author=AUTHOR_MODEL,
|
||||
@@ -782,7 +801,9 @@ async def _run_impression_write(context: ToolContext, args: dict[str, Any]) -> T
|
||||
if user is None:
|
||||
return _persona_error("impression_write", "There is nobody here to describe.")
|
||||
if not content:
|
||||
row = personas_service.impression(db, context.model_id, user)
|
||||
row = personas_service.impression(
|
||||
db, personas_service.key_for(context.model_id, context.data_group), user
|
||||
)
|
||||
if row is not None:
|
||||
personas_service.clear_impression(db, row)
|
||||
return ToolOutcome(
|
||||
@@ -791,7 +812,7 @@ async def _run_impression_write(context: ToolContext, args: dict[str, Any]) -> T
|
||||
)
|
||||
row = personas_service.write_impression(
|
||||
db,
|
||||
model_key=context.model_id,
|
||||
model_key=personas_service.key_for(context.model_id, context.data_group),
|
||||
owner=user,
|
||||
content=content,
|
||||
author=AUTHOR_MODEL,
|
||||
@@ -819,7 +840,7 @@ async def _run_memory_add(context: ToolContext, args: dict[str, Any]) -> ToolOut
|
||||
user = db.get(User, context.owner_id)
|
||||
try:
|
||||
memory = memories_service.add(
|
||||
db, owner=user, content=content, author=AUTHOR_MODEL
|
||||
db, owner=user, content=content, author=AUTHOR_MODEL, group=context.data_group
|
||||
)
|
||||
except ValueError as exc:
|
||||
return ToolOutcome(
|
||||
@@ -861,7 +882,7 @@ async def _run_memory_forget(context: ToolContext, args: dict[str, Any]) -> Tool
|
||||
wanted = str(args.get("content") or "").strip().lower()
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
records = memories_service.all_for(db, user)
|
||||
records = memories_service.all_for(db, user, context.data_group)
|
||||
if not wanted:
|
||||
return ToolOutcome(
|
||||
"Say which memory to remove, quoting its text.",
|
||||
@@ -918,6 +939,7 @@ async def _run_report_write(context: ToolContext, args: dict[str, Any]) -> ToolO
|
||||
source=SOURCE_CHAT,
|
||||
source_id=context.chat_id or "",
|
||||
model_id=context.model_id or "",
|
||||
group=context.data_group,
|
||||
)
|
||||
return ToolOutcome(
|
||||
f"Filed report {report.id} — {report.title!r}. "
|
||||
@@ -937,9 +959,11 @@ async def _run_report_search(context: ToolContext, args: dict[str, Any]) -> Tool
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
found = (
|
||||
reports_service.search(db, user, query, limit=8, vector=vector)
|
||||
reports_service.search(
|
||||
db, user, query, limit=8, vector=vector, group=context.data_group
|
||||
)
|
||||
if query
|
||||
else reports_service.recent(db, user, limit=8)
|
||||
else reports_service.recent(db, user, limit=8, group=context.data_group)
|
||||
)
|
||||
event = {
|
||||
"name": "report_search",
|
||||
@@ -962,7 +986,7 @@ async def _run_report_search(context: ToolContext, args: dict[str, Any]) -> Tool
|
||||
async def _run_report_get(context: ToolContext, args: dict[str, Any]) -> ToolOutcome:
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
report = reports_service.get(db, str(args.get("id") or ""), user)
|
||||
report = reports_service.get(db, str(args.get("id") or ""), user, context.data_group)
|
||||
if report is None:
|
||||
return ToolOutcome(
|
||||
"There is no such report.",
|
||||
@@ -985,7 +1009,7 @@ async def _run_skill_get(context: ToolContext, args: dict[str, Any]) -> ToolOutc
|
||||
name = str(args.get("name") or "").strip()
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
skill = skills_service.by_name(db, name, user)
|
||||
skill = skills_service.by_name(db, name, user, context.data_group)
|
||||
# Enforced here and not only in the listing. Without this the per-chat
|
||||
# narrowing is advisory: a model can name a skill it was never shown --
|
||||
# from an earlier turn, from a note -- and the runner would fetch it.
|
||||
@@ -1020,6 +1044,7 @@ async def _run_skill_create(context: ToolContext, args: dict[str, Any]) -> ToolO
|
||||
description=str(args.get("description") or ""),
|
||||
body=str(args.get("body") or ""),
|
||||
author=AUTHOR_MODEL,
|
||||
group=context.data_group,
|
||||
)
|
||||
except skills_service.SkillError as exc:
|
||||
return ToolOutcome(
|
||||
@@ -1039,7 +1064,9 @@ async def _run_skill_create(context: ToolContext, args: dict[str, Any]) -> ToolO
|
||||
async def _run_skill_edit(context: ToolContext, args: dict[str, Any]) -> ToolOutcome:
|
||||
with session_scope() as db:
|
||||
user = db.get(User, context.owner_id)
|
||||
skill = skills_service.by_name(db, str(args.get("name") or ""), user)
|
||||
skill = skills_service.by_name(
|
||||
db, str(args.get("name") or ""), user, context.data_group
|
||||
)
|
||||
if skill is None or skill.owner_id != context.owner_id:
|
||||
return ToolOutcome(
|
||||
"There is no such skill, or it belongs to someone else.",
|
||||
@@ -1878,7 +1905,16 @@ def resolve_tools(
|
||||
# something on here would still be reaching for a tool the gates had
|
||||
# already removed.
|
||||
off = scoped_off(chat)
|
||||
empty_library = not skills_service.count_enabled(db, user, exclude=scoped_skills_off(chat))
|
||||
# Counted in the answering model's own data group: skills in another group
|
||||
# are not readable here, so they must not keep `skill_get` on offer.
|
||||
from lembas.services import data_groups
|
||||
|
||||
empty_library = not skills_service.count_enabled(
|
||||
db,
|
||||
user,
|
||||
exclude=scoped_skills_off(chat),
|
||||
group=data_groups.for_speaker(db, user, chat, speaker),
|
||||
)
|
||||
|
||||
# What a crowd speaker may do, which is narrower than what the chat may.
|
||||
if crowd_turn is not None:
|
||||
@@ -2095,6 +2131,7 @@ def context_for(
|
||||
model.
|
||||
"""
|
||||
from lembas.services import chat as chat_service
|
||||
from lembas.services import data_groups
|
||||
from lembas.services.agent import session as agent_session
|
||||
|
||||
if chat is not None and speaker is None:
|
||||
@@ -2110,6 +2147,11 @@ def context_for(
|
||||
image_checkpoint=(chat.image_checkpoint or "") if chat is not None else "",
|
||||
model_id=(speaker.model_id or "") if speaker is not None else "",
|
||||
connection_id=(speaker.connection_id or "") if speaker is not None else "",
|
||||
data_group=(
|
||||
data_groups.for_speaker(db, user, chat, speaker)
|
||||
if chat is not None
|
||||
else DEFAULT_GROUP
|
||||
),
|
||||
base_ids=[base.id for base in chat.knowledge_bases] if chat is not None else [],
|
||||
skills_off=scoped_skills_off(chat),
|
||||
tools=tools.by_name if tools is not None else None,
|
||||
|
||||
Reference in New Issue
Block a user