Skip to content

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).

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.

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)

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)

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.

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.