msflib.conversation¶
Public modules of the msflib-conversation package.
msflib.conversation.actions
¶
Database access for the conversation domain.
All table access goes through these action classes; mutations use the
ModelAction super() methods so ModelBase lifecycle events
(conversation-create-pre-commit etc.) fire automatically.
ConversationAction is the only action that reads by
:class:~msflib.conversation.contracts.ConversationScope (tenant_id +
workspace_id) — it is the action behind resolve_conversation (§5.3),
turning a caller-supplied conversation_id string into a validated row.
Every other action in this module (members, threads, messages, reactions,
pins, resource links) operates on the already-resolved integer
conversation_id FK, not the scope contract — see architecture doc §5.2.
MemberAction also owns membership-lifecycle persistence
(activate_membership/deactivate_membership/change_role): the
membership-row mutation, its system_event timeline entry (via an
injected MessageAction), and its domain event, all in one transaction —
mirroring how ConversationAction owns conversation-row persistence.
services/membership.py keeps only policy and validation.
ConversationAction(*, settings: SettingsBase | None = None, member_action: MemberAction | None = None)
¶
Bases: ModelAction[Conversation, ConversationCreate, ConversationUpdate]
create_with_owner(session: Session, *, scope: ConversationScope, owner: Participant, data: ConversationCreate, commit: bool = True) -> tuple[Conversation, ConversationMember]
¶
Conversation row plus its owner membership, in one transaction.
data is the (already kind/visibility-resolved) request-boundary
ConversationCreate, passed through as-is -- tenant_id/
workspace_id/created_by_kind/created_by_id are context
supplied via update=. Returns the owner membership alongside the
conversation so callers (ChannelService) don't have to re-query
for a row this method just inserted.
The insert runs inside a SAVEPOINT so a slug collision (slug is
unique per tenant+workspace, see models/conversation.py) only
unwinds the failed insert, not any surrounding transaction; it
surfaces as ValueError (matching ChannelService.create_channel's
own input-validation errors) rather than an uncaught
IntegrityError.
create_with_members(session: Session, *, scope: ConversationScope, kind: ConversationKind, visibility: ConversationVisibility, created_by: Participant, members: Sequence[tuple[MemberKind, str]], direct_key: str | None = None, requesting: tuple[MemberKind, str] | None = None, commit: bool = True) -> tuple[Conversation, ConversationMember | None]
¶
Conversation row plus a flat set of member rows, in one transaction.
direct_key (dm/group_dm only) carries the unique constraint that
makes services.channels.create_direct's get-or-create race-safe —
see models/conversation.py's ux_conversation_scope_direct_key.
requesting (kind, id) picks out which of members' created
rows to also return -- so a caller who already knows which member is
"themselves" doesn't have to re-query for it (members has no
inherent order to key off of otherwise).
get_or_create_direct(session: Session, *, scope: ConversationScope, kind: ConversationKind, direct_key: str, created_by: Participant, members: Sequence[tuple[MemberKind, str]], requesting: tuple[MemberKind, str] | None = None) -> tuple[Conversation, ConversationMember | None]
¶
Get-or-create a dm/group_dm by its direct_key — a deterministic
hash of the member set (see services.channels._direct_key),
carrying the unique constraint in models/conversation.py's
ux_conversation_scope_direct_key. One closed unit of work
(guidelines §4): commits itself, no commit passthrough.
requesting (kind, id): see create_with_members; on the
get-existing path (no insert happens) this falls back to one
member_action.get_membership lookup, same cost as before.
Concurrency: two concurrent calls for the same member set can both miss the lookup below and both attempt to insert; the unique constraint lets only one through. The insert runs inside a SAVEPOINT so a losing IntegrityError only unwinds that insert (not any surrounding transaction); the loser then re-reads and returns the winner's row instead of raising or creating a duplicate.
list_visible(session: Session, *, scope: ConversationScope, participant: Participant, eligibility: MemberEligibility, kinds: list[ConversationKind] | None = None, include_archived: bool = False, offset: int = 0, limit: int = 50) -> list[tuple[Conversation, ConversationMember | None]]
¶
find_scoped plus the "mine or public" visibility filter and
the caller's own membership row, in one query — a conversation is
visible if participant is a member of it, or it's public and
eligibility admits participant (same policy can_read
applies to a single conversation, e.g. §6's workspace-membership
eligibility -- a caller a host's policy would reject from resolving
one public conversation must not see it in a list either).
eligibility.is_eligible is an arbitrary Python predicate, not a
SQL condition, so the public/non-member branch can't be filtered in
the WHERE clause. Rather than materializing every visible
conversation to filter+paginate in Python, the query re-runs with a
growing LIMIT (never OFFSET) until enough eligible rows are
collected to cover offset + limit -- each retry recomputes
collected from that single query's rows rather than stitching
together separate round trips, so a last_message_at bump (every
new message reorders the list) between retries can't cause rows to
be skipped or duplicated the way paging via OFFSET across
multiple queries would.
MemberAction(*, settings: SettingsBase | None = None, message_action: MessageAction | None = None)
¶
Bases: ModelAction[ConversationMember, ConversationMemberCreate, ConversationMemberUpdate]
activate_membership(session: Session, *, conversation: Conversation, member_kind: MemberKind, member_id: str, role: MemberRole, invited_by_id: str | None, action_label: str, event_key: str = 'member_joined', commit: bool = True) -> ConversationMember
¶
Create or reactivate a membership row, plus its timeline event and domain event, in one transaction. Idempotent: a no-op (no new timeline entry or domain event) if already active.
Concurrency: two concurrent joins/invites for the same never-before
member can both miss the existing check below and both attempt
to insert; ux_conversation_member_identity (the unique index in
models/member.py) lets only one insert through. The insert runs
inside a SAVEPOINT so a losing IntegrityError only unwinds that
insert (not the caller's surrounding transaction); the loser then
re-reads the winner's now-active row and returns it as a no-op,
rather than raising or double-emitting the join event.
deactivate_membership(session: Session, *, conversation: Conversation, membership: ConversationMember, new_status: MemberStatus, action_label: str, event_key: str, commit: bool = True) -> ConversationMember
¶
Leave/remove a membership row, plus its timeline event and domain event, in one transaction.
change_role(session: Session, *, conversation: Conversation, target_membership: ConversationMember, new_role: MemberRole, commit: bool = True) -> ConversationMember
¶
Change a member's role, plus its timeline event and domain event, in one transaction.
get_for_update(session: Session, *, member_id: int) -> ConversationMember
¶
Re-select the membership row with FOR UPDATE, so a caller
advancing last_read_message_id locks out concurrent advances
instead of both reading the same stale value and one commit
clobbering the other with a smaller id.
has_other_active_owner(session: Session, *, conversation_id: int, exclude_member_id: int) -> bool
¶
Existence check, not a bounded scan — correct regardless of how many members a conversation has.
Concurrency: locks every active-owner row for this conversation
(FOR UPDATE, in a fixed id order) rather than only the
"other" owners. Two owners leaving at the same instant both run this
query; without the lock both could see "another owner exists" and
both walk away, leaving zero owners. Locking the same row set in the
same order for every caller means the second transaction blocks
until the first commits, instead of each locking only the other's
row and deadlocking.
list_for_member(session: Session, *, scope: ConversationScope, member_kind: MemberKind, member_id: str, statuses: list[MemberStatus] | None = None, offset: int = 0, limit: int = 200) -> list[ConversationMember]
¶
List a member's memberships, scoped to a tenant+workspace.
ConversationMember carries no scope columns of its own —
conversation_id is the only source of truth, so both tenant and
workspace scoping join through Conversation. Without this join,
member ids colliding or being reused across tenants or workspaces
would leak memberships between them.
ThreadAction(*, settings: SettingsBase | None = None)
¶
Bases: ModelAction[ConversationThread, ConversationThreadCreate, ConversationThreadUpdate]
get_scoped(session: Session, *, thread_public_id: str, tenant_id: int, workspace_id: int | None) -> tuple[ConversationThread, Conversation] | None
¶
Thread + its conversation, tenant/workspace-filtered in one
joined query -- for AccessService.resolve_thread, so a
cross-tenant thread_public_id is rejected by the query itself
rather than fetched and checked in Python.
list_by_ids(session: Session, *, thread_ids: Sequence[int]) -> list[ConversationThread]
¶
Batch-fetch threads by id, in one query — used to resolve a page of messages' threads without an N+1 lookup per message.
get_for_update(session: Session, *, thread_id: int) -> ConversationThread
¶
Re-select the thread row with FOR UPDATE, so a caller
bumping reply_count locks out concurrent bumps instead of both
reading the same starting value and losing an increment.
get_or_create(session: Session, *, conversation: Conversation, root_message: Message, commit: bool = True) -> ConversationThread
¶
Get-or-create the thread row for a root message, plus its domain event, in one transaction.
Concurrency: two concurrent calls for the same root message can both
miss the lookup below and both attempt to insert; the unique
constraint on root_message_id lets only one through. The insert
runs inside a SAVEPOINT so a losing IntegrityError only unwinds that
insert (not any surrounding transaction); the loser then re-reads
and returns the winner's row instead of raising or creating a
duplicate.
MessageAction(*, settings: SettingsBase | None = None, conversation_action: ConversationAction | None = None, thread_action: ThreadAction | None = None)
¶
Bases: ModelAction[Message, MessageCreate, MessageUpdate]
post(session: Session, *, conversation: Conversation, sender_kind: SenderKind, sender_id: str, data: MessageCreate, thread: ConversationThread | None = None, message_type: MessageType = MessageType.message, mention_data: dict | None = None, extra_data: dict | None = None, commit: bool = True) -> Message
¶
Create a message, plus the conversation's last_message_at bump
and (for replies) the thread's reply-counter bump, in one
transaction. data is the request-boundary MessageCreate
(content/blocks/client_msg_id), passed through as-is to .create()
-- conversation_id/thread_id/sender_kind/sender_id/
message_type and any mention_data are context this method
resolves, supplied via update= rather than rebuilt onto a new
schema instance. extra_data is for trusted internal callers only
(e.g. the AI bridge tagging a turn with its agent_run_id) --
never the HTTP request boundary, which only ever supplies data
via MessageCreate; mention_data wins on key collision.
Idempotent on client_msg_id (scoped to
conversation + sender): a retry with the same id and the same
destination (thread, or main timeline) returns the original message
rather than creating a duplicate. Reusing an id for a different
destination (e.g. main timeline then a thread reply) is a client
bug, not a retry, and raises ValueError rather than silently
returning the wrong message.
Concurrency: two concurrent retries with the same client_msg_id
can both miss the lookup below and both attempt to insert; the
unique constraint on (conversation_id, sender_kind, sender_id,
client_msg_id) lets only one through. The insert runs inside a
SAVEPOINT so a losing IntegrityError only unwinds that insert (not
any surrounding transaction); the loser then re-reads and returns
the winner's row instead of raising or creating a duplicate. The
thread's reply_count bump locks the row (FOR UPDATE) first,
so concurrent replies don't read the same starting count and lose
an increment.
edit(session: Session, *, conversation: Conversation, message: Message, new_content: str, new_data: dict, commit: bool = True) -> Message
¶
Update a message's content, plus its domain event, in one
transaction. new_data is the full replacement for data
(edit-history bookkeeping is the caller's job). The event's
thread_id is resolved from message.thread_id rather than
accepted as a parameter, so it can't disagree with the message.
soft_delete(session: Session, *, conversation: Conversation, message: Message, commit: bool = True) -> Message
¶
Tombstone a message (deleted_at only; row/content untouched),
plus its domain event, in one transaction. The event's thread_id
is resolved from message.thread_id rather than accepted as a
parameter, so it can't disagree with the message.
get_scoped(session: Session, *, message_id: int, tenant_id: int, workspace_id: int | None) -> tuple[Message, Conversation] | None
¶
Message + its conversation, tenant/workspace-filtered in one
joined query -- for AccessService.resolve_message, so a
cross-tenant message_id is rejected by the query itself rather
than fetched and checked in Python.
list_for_conversation(session: Session, *, conversation_id: int, before_id: int | None = None, after_id: int | None = None, limit: int = 50) -> list[Message]
¶
Keyset pagination over the main timeline (thread_id IS NULL).
Newest-first when paginating with before_id (or no cursor);
after_id walks forward instead (e.g. "load newer"). Passing both
is not meaningful and before_id wins.
list_for_thread(session: Session, *, thread_id: int, before_id: int | None = None, limit: int = 50) -> list[Message]
¶
Keyset pagination over thread replies, oldest-first.
before_id walks backward from a cursor, so the page closest to
it is fetched newest-first (desc) and reversed for display —
ordering ascending under id < before_id would instead return
the oldest replies in the thread on every page.
list_by_ids(session: Session, *, message_ids: Sequence[int]) -> list[Message]
¶
Batch-fetch messages by id, in one query — used to resolve pinned messages without an N+1 lookup per pin.
count_after(session: Session, *, conversation_id: int, after_id: int, thread_id: int | None = None) -> int
¶
Count non-deleted main-timeline messages (or, if thread_id is
given, that thread's replies) posted after after_id — the
read_state.py unread-count query.
ReactionAction(*, settings: SettingsBase | None = None)
¶
Bases: ModelAction[Reaction, ReactionCreate, ReactionUpdate]
add(session: Session, *, conversation: Conversation, message_id: int, member_kind: MemberKind, member_id: str, emoji: str, commit: bool = True) -> Reaction
¶
Create a reaction, plus its domain event, in one transaction. Idempotent: a repeat of the same (message, member, emoji) is a no-op, not a duplicate row.
Concurrency: two concurrent adds of the same (message, member,
emoji) can both miss the lookup below and both attempt to insert;
ux_reaction_identity (the unique index in models/reaction.py)
lets only one through. The insert runs inside a SAVEPOINT so a
losing IntegrityError only unwinds that insert (not any surrounding
transaction); the loser then re-reads and returns the winner's row
instead of raising or creating a duplicate.
remove(session: Session, *, conversation: Conversation, message_id: int, member_kind: MemberKind, member_id: str, emoji: str, commit: bool = True) -> Reaction | None
¶
Delete a reaction, plus its domain event, in one transaction.
A no-op (no event, returns None) if the reaction doesn't exist.
PinAction(*, settings: SettingsBase | None = None)
¶
Bases: ModelAction[Pin, PinCreate, PinUpdate]
list_for_conversation(session: Session, *, conversation_id: int, offset: int = 0, limit: int | None = None) -> list[Pin]
¶
Ordered by PK for stable pin order (row order isn't otherwise guaranteed across databases/drivers).
add(session: Session, *, conversation_id: int, message_id: int, pinned_by_id: str, commit: bool = True) -> Pin
¶
Create a pin. Idempotent: a repeat pin of the same (conversation, message) is a no-op, not a duplicate row.
Concurrency: two concurrent pins of the same (conversation, message)
can both miss the lookup below and both attempt to insert;
ux_pin_conversation_message (the unique index in
models/pin.py) lets only one through. The insert runs inside a
SAVEPOINT so a losing IntegrityError only unwinds that insert (not
any surrounding transaction); the loser then re-reads and returns
the winner's row instead of raising or creating a duplicate.
ResourceLinkAction(*, settings: SettingsBase | None = None)
¶
Bases: ModelAction[ResourceLink, ResourceLinkCreate, ResourceLinkUpdate]
list_for_conversation(session: Session, *, conversation_id: int, offset: int = 0, limit: int = 100) -> list[ResourceLink]
¶
Ordered by PK for stable paging (row order isn't otherwise guaranteed across databases/drivers).
list_for_message(session: Session, *, conversation_id: int, message_id: int, offset: int = 0, limit: int | None = None) -> list[ResourceLink]
¶
Ordered by PK for stable paging (row order isn't otherwise guaranteed across databases/drivers).
link(session: Session, *, conversation: Conversation, data: ResourceLinkCreate, commit: bool = True) -> ResourceLink
¶
Create a resource link, plus its domain event, in one
transaction. conversation_id isn't a ResourceLinkCreate
field — it's always taken from conversation here.
msflib.conversation.config
¶
msflib.conversation.contracts
¶
Participant(kind: MemberKind, id: str, tenant_id: int)
dataclass
¶
The caller attempting an action, projected to membership's identity shape.
tenant_id must already be a real, resolved tenant id -- resolution
happens once at the request boundary (deps.py/ai_bridge.py),
matching workspace_id's existing (already-int) convention.
SameTenantEligibility
¶
Default eligibility: any participant in the conversation's tenant.
ConversationScope(tenant_id: int, workspace_id: int | None = None, conversation_id: str | None = None, sub_thread_id: str | None = None)
dataclass
¶
Validated scope filter values, typed to the conversation model columns.
workspace_id=None means the explicit workspace-less bucket (rendered
as IS NULL), matching the scope system's explicit-null semantics.
conversation_id/sub_thread_id identify a specific conversation
(and, optionally, thread within it) rather than acting as a hierarchical
visibility filter the way they do in the knowledge module — this module
is the one that resolves those strings into rows in the first place
(architecture doc §5.3), so most reads bind them by equality once a
conversation has been resolved. They are None for list-style
operations that are not yet scoped to one conversation (e.g. "list my
conversations").
identity
¶
Structural contracts for the account/workspace objects a host app passes
into deps.py -- mirrors msflib.account.contracts.WorkspaceContract's
own pattern (a local Protocol satisfied structurally, no inheritance
required) rather than importing it, so this module keeps zero hard
dependency on msflib-account/msflib-workspaces (architecture doc: no import
of account/workspaces/ai_core/documents at module scope).
participants
¶
Who is asking — the actor-side counterpart to :class:ConversationScope.
Participant is the boundary projection of an authenticated caller (human
account or agent) into the (kind, id-string) shape membership rows use
(architecture doc §3.1/§13.1) — never an FK into the account module.
MemberEligibility is the one seam access control delegates "who may
join a public channel / be invited" to (§6): the default is same-tenant;
the workspaces integration replaces it with workspace-membership checks
(§9.2) without this module ever importing workspaces.
Participant(kind: MemberKind, id: str, tenant_id: int)
dataclass
¶
The caller attempting an action, projected to membership's identity shape.
tenant_id must already be a real, resolved tenant id -- resolution
happens once at the request boundary (deps.py/ai_bridge.py),
matching workspace_id's existing (already-int) convention.
SameTenantEligibility
¶
Default eligibility: any participant in the conversation's tenant.
scope
¶
Typed compiled-scope contract for the conversation module.
Mirrors msflib.knowledge.contracts.KnowledgeScope: ScopeCompiler
validates a ScopeEnvelope against a profile and projects it to
string-typed sink formats; ConversationScope is the domain-typed
projection consumers of this module use — built once (by a canonical
per-module compiling helper, guidelines §8), then passed typed through
services -> actions with certain, non-defensive attribute access.
ConversationScope(tenant_id: int, workspace_id: int | None = None, conversation_id: str | None = None, sub_thread_id: str | None = None)
dataclass
¶
Validated scope filter values, typed to the conversation model columns.
workspace_id=None means the explicit workspace-less bucket (rendered
as IS NULL), matching the scope system's explicit-null semantics.
conversation_id/sub_thread_id identify a specific conversation
(and, optionally, thread within it) rather than acting as a hierarchical
visibility filter the way they do in the knowledge module — this module
is the one that resolves those strings into rows in the first place
(architecture doc §5.3), so most reads bind them by equality once a
conversation has been resolved. They are None for list-style
operations that are not yet scoped to one conversation (e.g. "list my
conversations").
transcript
¶
Transcript persistence contract (architecture doc §9.1.4).
The AI bridge writes agent turns as ordinary conversation messages
(sender_kind=agent) so the conversation becomes the human-readable
transcript of record, independent of whatever a LangGraph checkpointer keeps
as state. TranscriptWriter is the seam ai_api's router is handed as an
injected callable — ai_api never imports this module, it only calls whatever
object the host's factory returns (duck typing, same as the
ToolRegistry/get_tool_registries seam it already has).
msflib.conversation.deps
¶
Request-scoped FastAPI dependencies for router.py (architecture doc
§10). Mirrors ai_api.deps.get_ai_api_scope_dependencies's shape: a
factory that takes the host app's own get_current_account/
get_current_workspace callables and returns a DependencyNamespace of
dependencies that resolve a validated Participant, then a
ConversationContext/MessageContext/ThreadContext (row +
membership, access-checked via AccessService) -- so access checks
happen once, at the dependency layer, not repeated inline in every endpoint
body.
Every dependency here is overridable by a host app: build a namespace,
wrap or replace individual attributes, and pass the whole namespace back
into router(dependencies=...) -- same seam as workspaces.router's
workspace_form_dependencies.
Typed against this module's own AccountContract/WorkspaceContract
Protocols (only .id is required) rather than importing
account/workspaces model types, matching this module's
zero-hard-coupling design (architecture doc: no import of
account/ai_core/workspaces/documents at module scope).
get_conversation_dependencies(*, get_session: Callable, get_current_account: Callable, get_current_workspace: Callable | None = None, get_current_tenant: Callable | None = None, access_service: AccessService | None = None, conversation_action: ConversationAction | None = None, member_action: MemberAction | None = None, message_action: MessageAction | None = None, thread_action: ThreadAction | None = None) -> DependencyNamespace
¶
Build the conversation module's request-scoped dependency chain.
msflib.conversation.enums
¶
Conversation domain enums. Kept outside models so config-only callers
can import them without pulling in the ORM; models/enums.py re-exports
these for ORM-facing consumers.
msflib.conversation.eventbus
¶
events
¶
Conversation domain events (architecture doc §8).
Row-level lifecycle events (conversation-create-post-commit etc.) are
already emitted automatically by ModelAction — these are reserved for
aggregate/domain moments services emit explicitly. Payloads carry the
conversation's public_id (never the int PK) plus scope values so
listeners can act without re-querying; messages are the one entity this
module addresses externally by int id (see the HTTP surface, §10), so
message_id payload fields are plain ints.
msflib.conversation.integrations
¶
Optional integrations wiring msflib.conversation to other modules.
Each submodule pulls in one optional dependency (workspaces, ai_core/
ai_api, documents, ingestion, notifications, realtime)
and is imported directly by consumers that have that dependency installed —
never through msflib.conversation's lazy __getattr__, which must stay
free of optional imports. See architecture doc §9.
ai_bridge
¶
AI bridge — the integration layer between msflib.conversation and ai_api/ai_core.
Direction of coupling: this module exposes helpers that ai_api/host wiring
consumes; msflib.conversation never imports ai_core/ai_api at
module scope, and neither of those modules imports msflib.conversation.
The one place this file talks to ai_core (cleanup_ai_core_memories) does
so through a deferred import guarded by is_importable — the same
convention msflib.knowledge.integrations.worker uses for its own
optional ai_core call.
ConversationTranscriptWriter(make_session: Callable[[], Session], *, scope: ConversationScope, user: Participant, agent_id: str, agent_role: MemberRole = MemberRole.member, settings: ConversationSettings | None = None, eligibility: MemberEligibility | None = None)
¶
TranscriptWriter backed by services/messaging.py. Each write
opens its own short-lived session via make_session and re-resolves
the conversation/thread/memberships, rather than holding one session
open across a whole ask/stream request, so a slow or streamed agent turn
never holds a DB connection idle.
conversation_dimensions(conversation: Conversation, *, thread: ConversationThread | None = None) -> dict
¶
The validated {"conversation_id", "sub_thread_id"} scope dims for
ScopeEnvelope construction — resuming an agent becomes "posting in
the same channel/thread" for free once a caller builds its envelope from
this.
build_conversation_access_dependency(get_session: Callable[..., Any], get_current_account: Callable[..., Any], *, get_current_workspace: Callable[..., Any] | None = None, get_current_tenant: Callable[..., Any] | None = None, eligibility: MemberEligibility | None = None, settings: ConversationSettings | None = None) -> Callable[..., None]
¶
Build a FastAPI dependency that resolves + access-checks a request's
conversation_id before the wrapped route runs. Reads the request body
itself (conversation_id/sub_thread_id) via Request.json()
rather than declaring a typed Pydantic body model, since this module has
no ai_api AskRequest to import — it works with any JSON body shaped
like one. Wire it in front of ai_api's agent router without either module
importing the other::
api_router.include_router(
create_ai_router(...),
dependencies=[Depends(build_conversation_access_dependency(...))],
)
A request with no conversation_id in its body passes through
untouched. get_current_workspace mirrors deps.py's own dependency
chain — omitting it would 404 every workspace-scoped conversation, since
AccessService.resolve_conversation matches workspace_id exactly
(None means the explicit workspace-less bucket).
add_agent_participant(session: Session, *, conversation: Conversation, agent_id: str, role: MemberRole = MemberRole.member) -> ConversationMember
¶
Add an agent as a first-class conversation member (pass-through to
MembershipService.add_agent).
build_transcript_writer_factory(make_session: Callable[[], Session], *, agent_id: str, agent_role: MemberRole = MemberRole.member, settings: ConversationSettings | None = None, eligibility: MemberEligibility | None = None) -> Callable[..., TranscriptWriter | None]
¶
Per-request factory: transcript_writer_factory(scope, payload).
Returns None when the request carries no conversation_id.
build_speaker_label_factory(resolve_display_name: Callable[[str], str | None] | None = None) -> Callable[[Any, Any], str | None]
¶
Per-request factory, same shape as transcript_writer_factory.
cleanup_ai_core_memories(store: Any, *, tenant_id: int | None, conversation_id: str, memory_type: str, workspace_id: int | None = None) -> int
¶
Delete memory-store records namespaced to conversation_id — called
when a conversation's messages are purged, so episodic/semantic memory
doesn't outlive the transcript it was extracted from. No-ops (returns
0) if msflib.ai_core isn't installed.
MemoryStoreScope's aicore.memory profile has no
conversation_id dimension of its own, so conversation identity is
threaded through as sub_namespace instead (an opaque segment
scoped_namespace appends after memory_type).
register_memory_cleanup_listener(store: Any, *, memory_type: str, app: Any = None) -> Callable[..., Any] | None
¶
Subscribe cleanup_ai_core_memories to
conversation-messages-purged so a purge also clears the memory
namespace it just purged. No-ops if msflib.ai_core isn't installed.
Pass app= to bind to a specific app-scoped emitter instead of the
process-global one.
msflib.conversation.models
¶
ConversationCreate
¶
Bases: SchemaBase
Also the router's request body for creating any conversation kind --
tenant_id/workspace_id (request scope), created_by_kind/
created_by_id (the authenticated principal), and data are all
resolved server-side and supplied via update= in
ConversationAction.create_with_owner/create_with_members, never
accepted from a caller. visibility is optional -- omitting it means
"use this deployment's configured default", not "private"; see
ChannelService.create_channel's data.visibility or
self.settings.default_visibility fallback. participants is a
dm/group_dm's initial member set (ChannelService.create_direct); the
caller is added automatically, so only the other participants are
listed here -- ignored for channel/personal kinds.
ConversationRead
¶
Bases: SchemaBase
The HTTP-facing conversation shape (architecture doc §10).
ConversationUpdate
¶
Bases: SchemaBase
Request body for the one PATCH endpoint, which may rename/retopic
and change visibility in the same request -- ChannelService.update
validates both, then threads this straight through as
ConversationAction.update's data=, one commit. Deliberately
excludes slug and the internal state fields is_archived/
archived_at/last_message_at/data -- those are mutated by
other, dedicated endpoints/internal flows (ConversationAction.
set_archived, MessageAction.post's last_message_at bump) via
plain update= dicts, never through a schema instance a client could
supply values for.
MarkReadRequest
¶
Bases: SchemaBase
Mutates ConversationMemberFieldsMixin.last_read_message_id for the
caller's own membership -- narrower than ConversationMemberUpdate
for the same reason as MemberRoleChangeRequest.
MemberAddRequest
¶
Bases: SchemaBase
HTTP request body for inviting a member. No 1:1 ConversationMemberCreate
equivalent to reuse: member_kind is deliberately absent here (this
router only ever creates MemberKind.user members — see router.py's
module docstring), and conversation_id/status/invited_by_id
are resolved server-side, not accepted from the caller.
MemberRef
¶
Bases: SchemaBase
One entry of a dm/group_dm's initial member set — nested inside
ConversationCreate.participants, not a table row on its own.
MemberRoleChangeRequest
¶
Bases: SchemaBase
A deliberately narrower request than ConversationMemberUpdate,
which is shared by several distinct operations (role change, leave/
remove, mute, star, mark-read) -- this endpoint only ever accepts a
role.
MessageCreate
¶
Bases: SchemaBase
Also doubles as the router's request body for posting a message
(MessagePostRequest's former job) — everything else a row needs
(conversation_id, thread_id, sender_kind/sender_id,
message_type, mention data) is resolved server-side and supplied
via update= in MessageAction.post(), never accepted from a
caller.
MessageEditRequest
¶
Bases: SchemaBase
A deliberately narrower request than MessageUpdate, which is
shared by several distinct operations (edit, soft-delete, redact) --
this endpoint only ever accepts new content.
PinUpdate
¶
Bases: SchemaBase
Pins are add/remove-only; no mutable fields. Present only to satisfy
ModelAction's generic contract.
ReactionUpdate
¶
Bases: SchemaBase
Reactions are add/remove-only; no mutable fields. Present only to
satisfy ModelAction's generic contract.
ResourceLinkCreate
¶
Bases: SchemaBase
Also doubles as the router's request body -- conversation_id is
always resolved from the path (see ResourceLinkAction.link, which
supplies it via update=), so it isn't a field here at all.
conversation
¶
The conversation (channel/dm/group_dm/personal) — the container Slack
calls a channel. public_id is the conversation_id scope-dimension
value everywhere (checkpoints, memory namespaces, knowledge scope); never
expose or accept the int PK outside this module. See architecture doc §3.
ConversationCreate
¶
Bases: SchemaBase
Also the router's request body for creating any conversation kind --
tenant_id/workspace_id (request scope), created_by_kind/
created_by_id (the authenticated principal), and data are all
resolved server-side and supplied via update= in
ConversationAction.create_with_owner/create_with_members, never
accepted from a caller. visibility is optional -- omitting it means
"use this deployment's configured default", not "private"; see
ChannelService.create_channel's data.visibility or
self.settings.default_visibility fallback. participants is a
dm/group_dm's initial member set (ChannelService.create_direct); the
caller is added automatically, so only the other participants are
listed here -- ignored for channel/personal kinds.
ConversationRead
¶
Bases: SchemaBase
The HTTP-facing conversation shape (architecture doc §10).
ConversationUpdate
¶
Bases: SchemaBase
Request body for the one PATCH endpoint, which may rename/retopic
and change visibility in the same request -- ChannelService.update
validates both, then threads this straight through as
ConversationAction.update's data=, one commit. Deliberately
excludes slug and the internal state fields is_archived/
archived_at/last_message_at/data -- those are mutated by
other, dedicated endpoints/internal flows (ConversationAction.
set_archived, MessageAction.post's last_message_at bump) via
plain update= dicts, never through a schema instance a client could
supply values for.
enums
¶
Re-exports msflib.conversation.enums for existing .models.enums imports.
ids
¶
Public-id minting for scope-dimension-carrying rows.
Conversation.public_id / ConversationThread.public_id are not display
convenience — they ARE the conversation_id / sub_thread_id scope
dimension values (architecture doc §3.2/§5.3), so they must be non-enumerable
and mintable without coordination. ULIDs give lexicographic = chronological
ordering (useful for listing/debugging) without leaking a sequential count
the way an int PK would.
member
¶
Conversation membership — who belongs, with role/state and read state.
Members are polymorphic (member_kind + member_id) rather than FKs
into the account module: this is what makes "an AI participant in a channel"
a data statement instead of a special case, and keeps account/workspaces
optional dependencies (architecture doc §3.1/§13.1). Read state is embedded
here rather than a separate table — one row per member-conversation already
exists.
MemberRef
¶
Bases: SchemaBase
One entry of a dm/group_dm's initial member set — nested inside
ConversationCreate.participants, not a table row on its own.
MemberAddRequest
¶
Bases: SchemaBase
HTTP request body for inviting a member. No 1:1 ConversationMemberCreate
equivalent to reuse: member_kind is deliberately absent here (this
router only ever creates MemberKind.user members — see router.py's
module docstring), and conversation_id/status/invited_by_id
are resolved server-side, not accepted from the caller.
MemberRoleChangeRequest
¶
Bases: SchemaBase
A deliberately narrower request than ConversationMemberUpdate,
which is shared by several distinct operations (role change, leave/
remove, mute, star, mark-read) -- this endpoint only ever accepts a
role.
MarkReadRequest
¶
Bases: SchemaBase
Mutates ConversationMemberFieldsMixin.last_read_message_id for the
caller's own membership -- narrower than ConversationMemberUpdate
for the same reason as MemberRoleChangeRequest.
message
¶
Messages — the hot table. Ordering/pagination is keyset on
(conversation_id, id): the int PK is monotonic per insert order, so no
created_at tiebreak gymnastics are needed (architecture doc §3.2).
Soft delete (deleted_at) leaves a tombstone rather than removing the row,
so thread reply counts and knowledge retraction stay consistent; hard
deletion is the retention pass's job (§7 retention.py).
MessageCreate
¶
Bases: SchemaBase
Also doubles as the router's request body for posting a message
(MessagePostRequest's former job) — everything else a row needs
(conversation_id, thread_id, sender_kind/sender_id,
message_type, mention data) is resolved server-side and supplied
via update= in MessageAction.post(), never accepted from a
caller.
MessageEditRequest
¶
Bases: SchemaBase
A deliberately narrower request than MessageUpdate, which is
shared by several distinct operations (edit, soft-delete, redact) --
this endpoint only ever accepts new content.
pin
¶
PinUpdate
¶
Bases: SchemaBase
Pins are add/remove-only; no mutable fields. Present only to satisfy
ModelAction's generic contract.
reaction
¶
ReactionUpdate
¶
Bases: SchemaBase
Reactions are add/remove-only; no mutable fields. Present only to
satisfy ModelAction's generic contract.
resource_link
¶
Generic attachment/reference record — the decoupling seam for files and
anything else. Core stores and lists these; resolving a ref into a live
object is the owning integration's job (architecture doc §3.2/§9.3).
resource_kind is registry-validated by integrations, not enum-frozen, so
new kinds don't require a schema migration here.
ResourceLinkCreate
¶
Bases: SchemaBase
Also doubles as the router's request body -- conversation_id is
always resolved from the path (see ResourceLinkAction.link, which
supplies it via update=), so it isn't a field here at all.
thread
¶
First-class sub-conversation anchored to a root message (Slack-style
threads). First-class rather than a bare parent_message_id because the
sub_thread_id scope dimension and checkpoint-namespace hierarchy need a
stable identity to point at — see architecture doc §3.1/§13.2.
msflib.conversation.router
¶
HTTP surface (architecture doc §10). Mounted by the host app, which
injects its own session/account/workspace dependencies -- same factory
pattern as documents.router.router / workspaces.router.router, not a
fixed module-level router = APIRouter().
Three resource families share one router with no common path prefix applied
by default (prefix=""): /conversations/*, /messages/* (messages
are addressed externally by their own int id -- see eventbus/events.py --
not nested under a conversation path), and /threads/*. The host app can
still nest the whole thing under an outer prefix when mounting it.
HTTP participants are always MemberKind.user (the authenticated
account, stringified). Agent participants don't come through this router --
that's the (future) ai_bridge integration's job (§9.1).
Access checks (resolving a {public_id}/{message_id}/
{thread_public_id} path param into a row the caller may actually see)
happen once, at the dependency layer (deps.get_conversation_dependencies)
-- not repeated inline in every endpoint body. See deps.py for the
override seam (dependencies= here mirrors workspaces.router's
workspace_form_dependencies).
msflib.conversation.scope_profiles
¶
Scope profiles for the conversation module (architecture doc §5.1).
workspace_id is required_nullable (present, possibly None) rather
than merely optional (may be absent) — mirroring knowledge's
profiles — because a conversation's workspace-less-ness is a deliberate,
explicit fact (§3.2), not an unspecified one. conversation_id is a hard
required dimension for any operation bound to one conversation, since
this module is the one that resolves that string into a validated row
(§5.3); sub_thread_id stays optional since most operations address
the whole conversation, not one thread.
msflib.conversation.services
¶
AccessService(*, conversation_action: ConversationAction | None = None, member_action: MemberAction | None = None, message_action: MessageAction | None = None, thread_action: ThreadAction | None = None, eligibility: MemberEligibility | None = None)
¶
Wraps resolve_conversation/resolve_message/resolve_thread,
the functions in this module with action-class dependencies.
conversation_action/member_action/message_action/
thread_action/eligibility are injected once at construction —
they're wiring, not per-call data — while scope/participant
(genuinely per-request) stay method arguments.
resolve_conversation(session: Session, *, scope: ConversationScope, participant: Participant) -> tuple[Conversation, ConversationMember | None]
¶
Turn scope.conversation_id into a verified row (exists, tenant
matches, caller may read it), plus the membership row the access
check already had to load. Everything downstream — checkpointer
config, memory namespaces, knowledge scope — should flow through a
conversation resolved this way (§5.3).
resolve_message(session: Session, *, message_id: int, participant: Participant, workspace_id: int | None = None) -> tuple[Message, Conversation, ConversationMember | None]
¶
Turn a message id into a verified (message, conversation,
membership) triple — same read-access guarantee as
resolve_conversation, since a message is only visible if its
conversation is.
resolve_thread(session: Session, *, thread_public_id: str, participant: Participant, workspace_id: int | None = None) -> tuple[ConversationThread, Conversation, ConversationMember | None]
¶
Turn a thread's public_id into a verified (thread, conversation,
membership) triple — same read-access guarantee as
resolve_conversation, since a thread is only visible if its
conversation is.
ChannelService(*, conversation_action: ConversationAction | None = None, conversation_settings: ConversationSettings | None = None)
¶
create_channel(session: Session, *, scope: ConversationScope, data: ConversationCreate, created_by: Participant) -> tuple[Conversation, ConversationMember]
¶
data is the request-boundary ConversationCreate, passed
through to ConversationAction.create_with_owner unchanged apart
from resolving kind/visibility -- model_copy derives a
new instance rather than rebuilding one field-by-field. Returns the
owner membership alongside the conversation, mirroring
AccessService.resolve_conversation's (row, membership) shape,
so callers don't have to re-query for a row just inserted.
create_personal(session: Session, *, scope: ConversationScope, owner: Participant, data: ConversationCreate | None = None) -> tuple[Conversation, ConversationMember]
¶
A personal conversation ignores everything on data but
name -- kind/visibility/slug/topic/purpose
are always this method's own fixed values, not the caller's.
create_direct(session: Session, *, scope: ConversationScope, participants: Sequence[Participant], requested_kind: ConversationKind | None = None) -> tuple[Conversation, ConversationMember | None]
¶
Get-or-create a dm (2 participants) or group_dm (3+); no ad-hoc invite afterwards — starting a direct conversation with a different member set always yields a distinct conversation.
Race-safe get-or-create (concurrent requests for the same member set)
lives on ConversationAction.get_or_create_direct — this method
only computes the member set's direct_key and derives kind.
The returned membership is participants[0]'s -- the requesting
caller's, by this method's own convention (see router.py's
create_conversation) -- or None if it somehow isn't in the
deduped member set. requested_kind, if given, must match the kind
derived from the deduped participant count -- a caller asking for a
dm with 3+ people (or vice versa) is a client error, not silently
the other kind.
update(session: Session, *, conversation: Conversation, actor_membership: ConversationMember | None, data: ConversationUpdate) -> Conversation
¶
data is the request-boundary ConversationUpdate for
the one PATCH endpoint, which may rename/retopic and change
visibility in the same request -- both are validated before either
is written, then data is threaded through unchanged as
ConversationAction.update's data= (its own exclude_unset
dump already picks out only the fields the caller actually set), in
one commit -- a disallowed visibility transition can't leave an
allowed rename committed on its own this way. can_manage is
required unconditionally, even for a no-op request with no fields
set -- this is a management endpoint, not a public read. dm/group_dm
are explicitly rejected rather than left to can_manage's incidental
always-False result for their fixed MemberRole.member membership.
HistoryService(*, message_action: MessageAction | None = None, pin_action: PinAction | None = None, thread_action: ThreadAction | None = None, conversation_settings: ConversationSettings | None = None)
¶
list_conversation_history(session: Session, *, conversation: Conversation, before_id: int | None = None, after_id: int | None = None, limit: int | None = None) -> list[MessageRead]
¶
Main-timeline (thread_id IS NULL) keyset page, newest-first
unless after_id is given (see MessageAction.list_for_conversation).
around_message(session: Session, *, conversation: Conversation, message_id: int, before: int | None = None, after: int | None = None) -> list[MessageRead]
¶
A context window around message_id on the conversation's main
timeline, chronological order (oldest first). Use list_thread_history
for context inside a thread — a threaded message isn't on the main
timeline this method queries.
MembershipService(*, member_action: MemberAction | None = None, eligibility: MemberEligibility | None = None)
¶
join(session: Session, *, conversation: Conversation, participant: Participant) -> ConversationMember
¶
Self-serve join. Only valid for public, non-direct conversations.
invite(session: Session, *, conversation: Conversation, inviter_membership: ConversationMember | None, invitee: Participant, role: MemberRole = MemberRole.member) -> ConversationMember
¶
Any active, non-guest member may invite into a channel; DM/group_dm/ personal conversations don't support ad-hoc invite — start a new direct conversation with the desired member set instead (§14 Q2).
Granting an elevated role (admin/owner) requires the inviter to already
be an admin/owner (can_manage) — otherwise a plain member could
hand out ownership/admin to an arbitrary invitee, bypassing
change_role's "only the current owner may transfer ownership" gate.
The invitee is subject to the same MemberEligibility seam as
self-serve join (default: same-tenant) — this is what prevents
inviting a member_id belonging to another tenant, since membership rows
carry no tenant column of their own to enforce this at the DB layer.
leave(session: Session, *, conversation: Conversation, participant: Participant, membership: ConversationMember | None = None) -> ConversationMember
¶
DM/group_dm/personal membership is fixed at creation (§14 Q2):
leaving one is a dead end, since join/invite both reject
direct kinds and ChannelService.create_direct returns the same
conversation for the same member set regardless of who has left
it -- so a participant who leaves could never regain access.
membership, if given, must be participant's own row in
conversation -- lets a caller that already resolved it (e.g. the
router's access-check dependency) skip a redundant re-query, but a
mismatched row (wrong participant or wrong conversation) is a caller
bug, not silently trusted; otherwise it's looked up here.
remove(session: Session, *, conversation: Conversation, remover_membership: ConversationMember, target_kind: MemberKind, target_id: str) -> ConversationMember
¶
Same immutable-membership guard as leave -- removing someone
from a direct conversation is just as much a dead end as them
leaving it themselves.
add_agent(session: Session, *, conversation: Conversation, agent_id: str, role: MemberRole = MemberRole.member) -> ConversationMember
¶
Activate an agent membership row with no inviter to check — the AI
bridge integration wires a bot into a conversation programmatically
(architecture doc §9.1.3), unlike invite which is always a human
member's action reachable through the HTTP surface.
MessagingService(*, message_action: MessageAction | None = None, thread_action: ThreadAction | None = None, conversation_settings: ConversationSettings | None = None)
¶
post_message(session: Session, *, conversation: Conversation, sender: Participant, actor_membership: ConversationMember | None, data: MessageCreate, root_message: Message | None = None, extra_data: dict | None = None) -> Message
¶
data passes through to MessageAction.post unchanged; only mention_data
is added alongside it. With root_message, the thread is flushed (commit=False)
only after permission/content validation passes -- rejected requests then never touch
the DB for the thread -- then committed together with the reply; a rejected reply
rolls the thread back too. extra_data is for trusted internal callers only (see
MessageAction.post) -- never wired to the HTTP request boundary.
delete_message(session: Session, *, conversation: Conversation, message: Message, actor: Participant, actor_membership: ConversationMember | None) -> Message
¶
Soft delete: sender or admin/owner. Leaves a tombstone (row + FKs intact) so thread reply counts stay consistent; history.py redacts content/blocks for deleted rows on read.
ReadStateService(*, member_action: MemberAction | None = None, message_action: MessageAction | None = None)
¶
mark_read(session: Session, *, membership: ConversationMember, up_to_message_id: int) -> ConversationMember
¶
Advance last_read_message_id; never moves it backwards (an
out-of-order or stale client call is a silent no-op, not an error).
Locks the membership row (FOR UPDATE) before the read-then-write,
so two concurrent calls can't both read the same stale
last_read_message_id and commit out of order, leaving the marker
behind a value it already advanced past. Re-checks active status on
the locked row too (not just the caller-supplied membership), so
a removal that lands between the caller's read and this call still
blocks the advance.
mention_count(session: Session, *, participant: Participant, membership: ConversationMember, scan_limit: int = 500) -> int
¶
Count unread messages that @mention participant.id.
Scoped to the main timeline, like unread_count() — a mention
inside a thread reply isn't counted here; use list_thread_history
to inspect a specific thread's messages instead.
Scans up to scan_limit unread messages and filters in Python rather
than with a SQL-level JSON query — Message.data is a JSON blob and
the unread window is expected to be bounded, so a dialect-portable
(sqlite + Postgres) filter beats a dialect-specific JSON operator here.
ResourceKindRegistry(*, kinds: list[str] | None = None)
¶
A small, mutable set of known resource_kind strings.
ResourceService(*, resource_link_action: ResourceLinkAction | None = None, message_action: MessageAction | None = None, registry: ResourceKindRegistry | None = None)
¶
link_resource(session: Session, *, conversation: Conversation, actor_membership: ConversationMember | None, data: ResourceLinkCreate) -> ResourceLink
¶
data is the request-boundary ResourceLinkCreate, passed
through to ResourceLinkAction.link unchanged.
update_resource(session: Session, *, resource: ResourceLink, actor_membership: ConversationMember | None, data: ResourceLinkUpdate) -> ResourceLink
¶
ModelAction.update() already extracts only the fields the
caller actually set on data (exclude_unset) -- the guard
below just avoids an empty no-op write, not a manual field-by-field
dict build.
ThreadService(*, thread_action: ThreadAction | None = None)
¶
record_reply(session: Session, *, thread: ConversationThread) -> ConversationThread
¶
Locks the thread row (FOR UPDATE) before bumping
reply_count, so concurrent replies don't both read the same
starting count and lose an increment.
can_change_visibility(current: ConversationVisibility, new: ConversationVisibility, *, is_admin: bool, allow_private_to_public: bool) -> bool
¶
public -> private always allowed (to an admin); private -> public needs an admin and the config flag; personal never transitions. Admin is required even for a no-op (current == new) request — this is a permission gate, not just a state-transition check, so a non-admin caller must not get a successful response merely by requesting no change.
can_manage(conversation: Conversation, membership: ConversationMember | None) -> bool
¶
Rename, archive, member admin: role >= admin.
can_post(conversation: Conversation, membership: ConversationMember | None, *, allow_guest_posting: bool) -> bool
¶
Active member with role >= member; archived conversations are read-only.
can_read(conversation: Conversation, membership: ConversationMember | None, participant: Participant, *, eligibility: MemberEligibility | None = None) -> bool
¶
personal: owner only. private: active member only. public: active member, or an eligible non-member (read-before-join is allowed, Slack-style).
access
¶
Access control (architecture doc §6): pure policy over already-loaded
rows, plus AccessService.resolve_conversation (§5.3) — the single entry
point that turns a scope's conversation_id string into a verified row.
resolve_conversation does issue DB reads (it has to, to load the row and
membership it then checks), but performs no writes — it composes two action
reads with the pure can_read policy below, which is what keeps it here
rather than in channels.py/membership.py.
can_read/can_post/can_manage/can_change_visibility/
is_active_member/actor_matches_membership stay module-level
functions rather than methods: they take no action/settings dependencies,
only already-loaded rows and pure config values, so there is nothing to
inject at construction.
can_read/can_post/can_manage/is_active_member all check
that membership belongs to the conversation being evaluated, as part of
the access decision (a membership row from a different conversation should
never grant anything here) — kept as one check in one place rather than
duplicated at every call site.
AccessService(*, conversation_action: ConversationAction | None = None, member_action: MemberAction | None = None, message_action: MessageAction | None = None, thread_action: ThreadAction | None = None, eligibility: MemberEligibility | None = None)
¶
Wraps resolve_conversation/resolve_message/resolve_thread,
the functions in this module with action-class dependencies.
conversation_action/member_action/message_action/
thread_action/eligibility are injected once at construction —
they're wiring, not per-call data — while scope/participant
(genuinely per-request) stay method arguments.
resolve_conversation(session: Session, *, scope: ConversationScope, participant: Participant) -> tuple[Conversation, ConversationMember | None]
¶
Turn scope.conversation_id into a verified row (exists, tenant
matches, caller may read it), plus the membership row the access
check already had to load. Everything downstream — checkpointer
config, memory namespaces, knowledge scope — should flow through a
conversation resolved this way (§5.3).
resolve_message(session: Session, *, message_id: int, participant: Participant, workspace_id: int | None = None) -> tuple[Message, Conversation, ConversationMember | None]
¶
Turn a message id into a verified (message, conversation,
membership) triple — same read-access guarantee as
resolve_conversation, since a message is only visible if its
conversation is.
resolve_thread(session: Session, *, thread_public_id: str, participant: Participant, workspace_id: int | None = None) -> tuple[ConversationThread, Conversation, ConversationMember | None]
¶
Turn a thread's public_id into a verified (thread, conversation,
membership) triple — same read-access guarantee as
resolve_conversation, since a thread is only visible if its
conversation is.
can_read(conversation: Conversation, membership: ConversationMember | None, participant: Participant, *, eligibility: MemberEligibility | None = None) -> bool
¶
personal: owner only. private: active member only. public: active member, or an eligible non-member (read-before-join is allowed, Slack-style).
can_post(conversation: Conversation, membership: ConversationMember | None, *, allow_guest_posting: bool) -> bool
¶
Active member with role >= member; archived conversations are read-only.
can_manage(conversation: Conversation, membership: ConversationMember | None) -> bool
¶
Rename, archive, member admin: role >= admin.
is_active_member(membership: ConversationMember | None, conversation_id: int | None = None) -> bool
¶
Active-membership check. If conversation_id is given, also
requires membership to belong to that conversation — omit it only
when the caller has no separate conversation to check against (e.g.
membership.conversation_id is itself the scope, as in
read_state.py).
actor_matches_membership(membership: ConversationMember | None, participant: Participant) -> bool
¶
Does membership actually represent participant?
can_change_visibility(current: ConversationVisibility, new: ConversationVisibility, *, is_admin: bool, allow_private_to_public: bool) -> bool
¶
public -> private always allowed (to an admin); private -> public needs an admin and the config flag; personal never transitions. Admin is required even for a no-op (current == new) request — this is a permission gate, not just a state-transition check, so a non-admin caller must not get a successful response merely by requesting no change.
channels
¶
Conversation lifecycle (architecture doc §7): create, rename/topic/purpose, archive/unarchive, visibility transitions.
DM/group-DM creation is get-or-create on the normalised member set — a second "DM with the same people" returns the existing conversation rather than creating a duplicate (§7, §14 Q2: DM membership is immutable after creation; add someone new by starting a new direct conversation instead).
Row-level persistence (conversation + member rows, one transaction, plus the
conversation-created event) lives on ConversationAction
(create_with_owner/create_with_members) — ChannelService owns
defaults, validation, and the direct-conversation get-or-create policy, not
the writes.
ChannelService takes its conversation_action/conversation_settings
dependencies at construction, not per call. Every method is one closed unit
of work: it commits itself and takes no commit passthrough (guidelines
§4). Callers who need to compose several actions in one transaction call the
action classes directly.
ChannelService(*, conversation_action: ConversationAction | None = None, conversation_settings: ConversationSettings | None = None)
¶
create_channel(session: Session, *, scope: ConversationScope, data: ConversationCreate, created_by: Participant) -> tuple[Conversation, ConversationMember]
¶
data is the request-boundary ConversationCreate, passed
through to ConversationAction.create_with_owner unchanged apart
from resolving kind/visibility -- model_copy derives a
new instance rather than rebuilding one field-by-field. Returns the
owner membership alongside the conversation, mirroring
AccessService.resolve_conversation's (row, membership) shape,
so callers don't have to re-query for a row just inserted.
create_personal(session: Session, *, scope: ConversationScope, owner: Participant, data: ConversationCreate | None = None) -> tuple[Conversation, ConversationMember]
¶
A personal conversation ignores everything on data but
name -- kind/visibility/slug/topic/purpose
are always this method's own fixed values, not the caller's.
create_direct(session: Session, *, scope: ConversationScope, participants: Sequence[Participant], requested_kind: ConversationKind | None = None) -> tuple[Conversation, ConversationMember | None]
¶
Get-or-create a dm (2 participants) or group_dm (3+); no ad-hoc invite afterwards — starting a direct conversation with a different member set always yields a distinct conversation.
Race-safe get-or-create (concurrent requests for the same member set)
lives on ConversationAction.get_or_create_direct — this method
only computes the member set's direct_key and derives kind.
The returned membership is participants[0]'s -- the requesting
caller's, by this method's own convention (see router.py's
create_conversation) -- or None if it somehow isn't in the
deduped member set. requested_kind, if given, must match the kind
derived from the deduped participant count -- a caller asking for a
dm with 3+ people (or vice versa) is a client error, not silently
the other kind.
update(session: Session, *, conversation: Conversation, actor_membership: ConversationMember | None, data: ConversationUpdate) -> Conversation
¶
data is the request-boundary ConversationUpdate for
the one PATCH endpoint, which may rename/retopic and change
visibility in the same request -- both are validated before either
is written, then data is threaded through unchanged as
ConversationAction.update's data= (its own exclude_unset
dump already picks out only the fields the caller actually set), in
one commit -- a disallowed visibility transition can't leave an
allowed rename committed on its own this way. can_manage is
required unconditionally, even for a no-op request with no fields
set -- this is a management endpoint, not a public read. dm/group_dm
are explicitly rejected rather than left to can_manage's incidental
always-False result for their fixed MemberRole.member membership.
history
¶
Read paths (architecture doc §7): keyset-paginated channel/thread
timelines, an "around a message" context window, and pinned-message
listing. Soft-deleted messages stay in the result set as tombstones — the
row (and thread reply counts) stay intact — but content/blocks/
data are redacted here rather than at delete time, so the DB row itself
is never mutated beyond deleted_at.
message_action/pin_action/conversation_settings are injected at
construction, not per call.
HistoryService(*, message_action: MessageAction | None = None, pin_action: PinAction | None = None, thread_action: ThreadAction | None = None, conversation_settings: ConversationSettings | None = None)
¶
list_conversation_history(session: Session, *, conversation: Conversation, before_id: int | None = None, after_id: int | None = None, limit: int | None = None) -> list[MessageRead]
¶
Main-timeline (thread_id IS NULL) keyset page, newest-first
unless after_id is given (see MessageAction.list_for_conversation).
around_message(session: Session, *, conversation: Conversation, message_id: int, before: int | None = None, after: int | None = None) -> list[MessageRead]
¶
A context window around message_id on the conversation's main
timeline, chronological order (oldest first). Use list_thread_history
for context inside a thread — a threaded message isn't on the main
timeline this method queries.
present_message(message: Message, *, conversation: Conversation, thread: ConversationThread | None = None, redact: bool = True) -> MessageRead
¶
Redact a soft-deleted message's content/blocks/data and substitute
the resolved conversation/thread's public_id for their int FKs.
thread lets a caller skip message.thread's lazy load when
already in hand. redact=False (router.py's delete_message response)
lets the actor who just deleted a message see what they deleted.
membership
¶
Membership (architecture doc §7): join, invite, leave, remove, role
changes. MembershipService owns policy and validation only — the
membership-row mutation, its system_event timeline entry, and its domain
event are one transaction living on MemberAction
(activate_membership/deactivate_membership/change_role),
mirroring how ConversationAction owns conversation-row persistence (see
channels.py).
member_action/eligibility are injected at construction, not per
call. Every method is one closed unit of work: it commits itself and takes
no commit passthrough (guidelines §4). Callers who need to compose
several actions in one transaction call the action classes directly.
MembershipService(*, member_action: MemberAction | None = None, eligibility: MemberEligibility | None = None)
¶
join(session: Session, *, conversation: Conversation, participant: Participant) -> ConversationMember
¶
Self-serve join. Only valid for public, non-direct conversations.
invite(session: Session, *, conversation: Conversation, inviter_membership: ConversationMember | None, invitee: Participant, role: MemberRole = MemberRole.member) -> ConversationMember
¶
Any active, non-guest member may invite into a channel; DM/group_dm/ personal conversations don't support ad-hoc invite — start a new direct conversation with the desired member set instead (§14 Q2).
Granting an elevated role (admin/owner) requires the inviter to already
be an admin/owner (can_manage) — otherwise a plain member could
hand out ownership/admin to an arbitrary invitee, bypassing
change_role's "only the current owner may transfer ownership" gate.
The invitee is subject to the same MemberEligibility seam as
self-serve join (default: same-tenant) — this is what prevents
inviting a member_id belonging to another tenant, since membership rows
carry no tenant column of their own to enforce this at the DB layer.
leave(session: Session, *, conversation: Conversation, participant: Participant, membership: ConversationMember | None = None) -> ConversationMember
¶
DM/group_dm/personal membership is fixed at creation (§14 Q2):
leaving one is a dead end, since join/invite both reject
direct kinds and ChannelService.create_direct returns the same
conversation for the same member set regardless of who has left
it -- so a participant who leaves could never regain access.
membership, if given, must be participant's own row in
conversation -- lets a caller that already resolved it (e.g. the
router's access-check dependency) skip a redundant re-query, but a
mismatched row (wrong participant or wrong conversation) is a caller
bug, not silently trusted; otherwise it's looked up here.
remove(session: Session, *, conversation: Conversation, remover_membership: ConversationMember, target_kind: MemberKind, target_id: str) -> ConversationMember
¶
Same immutable-membership guard as leave -- removing someone
from a direct conversation is just as much a dead end as them
leaving it themselves.
add_agent(session: Session, *, conversation: Conversation, agent_id: str, role: MemberRole = MemberRole.member) -> ConversationMember
¶
Activate an agent membership row with no inviter to check — the AI
bridge integration wires a bot into a conversation programmatically
(architecture doc §9.1.3), unlike invite which is always a human
member's action reachable through the HTTP surface.
messaging
¶
Message posting/editing/deletion (architecture doc §7): access + size
validation, @mention extraction into Message.data["mentions"], and
client_msg_id idempotency. The actual write — a posted message, the
conversation's last_message_at bump, and (for replies) the thread's
reply-counter bump, all in one commit — lives on MessageAction.post
(mirroring how MemberAction owns membership-row persistence);
MessagingService owns policy and validation only.
message_action/thread_action/conversation_settings are injected
at construction, not per call. Every method is one closed unit of work: it
commits itself and takes no commit passthrough (guidelines §4).
post_message is the one exception -- it composes a thread get-or-create
with the post itself and rolls the session back on any failure, so neither
an orphaned thread row nor a half-flushed session survives a rejected
request. Callers who need to compose several actions in one transaction
elsewhere call the action classes directly.
MessagingService(*, message_action: MessageAction | None = None, thread_action: ThreadAction | None = None, conversation_settings: ConversationSettings | None = None)
¶
post_message(session: Session, *, conversation: Conversation, sender: Participant, actor_membership: ConversationMember | None, data: MessageCreate, root_message: Message | None = None, extra_data: dict | None = None) -> Message
¶
data passes through to MessageAction.post unchanged; only mention_data
is added alongside it. With root_message, the thread is flushed (commit=False)
only after permission/content validation passes -- rejected requests then never touch
the DB for the thread -- then committed together with the reply; a rejected reply
rolls the thread back too. extra_data is for trusted internal callers only (see
MessageAction.post) -- never wired to the HTTP request boundary.
delete_message(session: Session, *, conversation: Conversation, message: Message, actor: Participant, actor_membership: ConversationMember | None) -> Message
¶
Soft delete: sender or admin/owner. Leaves a tombstone (row + FKs intact) so thread reply counts stay consistent; history.py redacts content/blocks for deleted rows on read.
reactions
¶
Reactions and pins (architecture doc §7): both are add/remove-only
satellite tables keyed to a message, so both live on this one
ReactionService rather than splitting into two classes for a handful of
methods each. Reaction add/remove (row + domain event, one commit) lives on
ReactionAction; pins have no domain event so their methods delegate
straight to a single PinAction call.
reaction_action/pin_action are injected at construction, not per
call. Every method is one closed unit of work: it commits itself and takes
no commit passthrough (guidelines §4). Callers who need to compose
several actions in one transaction call the action classes directly.
read_state
¶
Read state (architecture doc §7): the monotonic last_read_message_id
marker embedded on ConversationMember (§3.2 — one row per member already
exists, so read state lives there rather than in its own table), plus
unread/mention counts for channel lists.
member_action/message_action are injected at construction, not per
call. mark_read is one closed unit of work: it commits itself and takes
no commit passthrough (guidelines §4). Callers who need to compose
several actions in one transaction call the action classes directly.
ReadStateService(*, member_action: MemberAction | None = None, message_action: MessageAction | None = None)
¶
mark_read(session: Session, *, membership: ConversationMember, up_to_message_id: int) -> ConversationMember
¶
Advance last_read_message_id; never moves it backwards (an
out-of-order or stale client call is a silent no-op, not an error).
Locks the membership row (FOR UPDATE) before the read-then-write,
so two concurrent calls can't both read the same stale
last_read_message_id and commit out of order, leaving the marker
behind a value it already advanced past. Re-checks active status on
the locked row too (not just the caller-supplied membership), so
a removal that lands between the caller's read and this call still
blocks the advance.
mention_count(session: Session, *, participant: Participant, membership: ConversationMember, scan_limit: int = 500) -> int
¶
Count unread messages that @mention participant.id.
Scoped to the main timeline, like unread_count() — a mention
inside a thread reply isn't counted here; use list_thread_history
to inspect a specific thread's messages instead.
Scans up to scan_limit unread messages and filters in Python rather
than with a SQL-level JSON query — Message.data is a JSON blob and
the unread window is expected to be bounded, so a dialect-portable
(sqlite + Postgres) filter beats a dialect-specific JSON operator here.
resources
¶
Resource links (architecture doc §3.2/§7/§9.3): the generic
attachment/reference seam. resource_kind is registry-validated rather
than enum-frozen so integrations can add kinds without a schema migration
here — Phase 2 owns storage/listing only; resolving a ref into a live
object is a later integration's job (§9.3), so no resolvers exist yet.
Linking (row + domain event, one commit) lives on
ResourceLinkAction.link.
resource_link_action/registry are injected at construction, not per
call. Every method is one closed unit of work: it commits itself and takes
no commit passthrough (guidelines §4). Callers who need to compose
several actions in one transaction call the action classes directly.
ResourceKindRegistry(*, kinds: list[str] | None = None)
¶
A small, mutable set of known resource_kind strings.
ResourceService(*, resource_link_action: ResourceLinkAction | None = None, message_action: MessageAction | None = None, registry: ResourceKindRegistry | None = None)
¶
link_resource(session: Session, *, conversation: Conversation, actor_membership: ConversationMember | None, data: ResourceLinkCreate) -> ResourceLink
¶
data is the request-boundary ResourceLinkCreate, passed
through to ResourceLinkAction.link unchanged.
update_resource(session: Session, *, resource: ResourceLink, actor_membership: ConversationMember | None, data: ResourceLinkUpdate) -> ResourceLink
¶
ModelAction.update() already extracts only the fields the
caller actually set on data (exclude_unset) -- the guard
below just avoids an empty no-op write, not a manual field-by-field
dict build.
threads
¶
Thread lifecycle (architecture doc §7): get-or-create the thread row for a
root message, and maintain the reply_count/last_reply_at counters
that MessagingService.post_message bumps on every threaded reply (via
MessageAction.post, which folds the counter bump into the same commit as
the reply itself). record_reply here is exposed for direct/standalone
use — normal reply flow does not call it. This class owns only the thread
row, not message posting.
thread_action is injected at construction, not per call. Every method is
one closed unit of work: it commits itself and takes no commit
passthrough (guidelines §4). Callers who need to compose several actions in
one transaction call the action classes directly.