"""Business logic for the kommunikation plugin.""" from __future__ import annotations import logging import re import uuid from datetime import datetime, timezone from typing import Any from sqlalchemy import select, update, func, and_, or_ from sqlalchemy.ext.asyncio import AsyncSession from app.core.event_bus import get_event_bus from app.plugins.builtins.kommunikation.models import ( CommConversation, CommConversationMute, CommConversationPin, CommMessage, CommMessageAttachment, CommMessageBlock, CommMessageEdit, CommMessageRead, CommMessageReaction, CommParticipant, ) from app.plugins.builtins.kommunikation.participant_registry import get_participant_registry from app.plugins.builtins.kommunikation.rbac import CommRBAC logger = logging.getLogger(__name__) MAX_TRIGGER_DEPTH = 3 # ─── Mention Parsing ─── MENTION_PATTERN = re.compile(r"@(\w+)") def parse_mentions(content: str) -> list[str]: """Parse @mentions from message content. Returns list of mentioned types/names.""" return MENTION_PATTERN.findall(content) # ─── Conversation Helpers ─── def conversation_to_response( conv: CommConversation, participants: list[CommParticipant], unread_count: int = 0, is_pinned_by_user: bool = False, ) -> dict[str, Any]: """Convert a CommConversation to a response dict.""" return { "id": str(conv.id), "title": conv.title, "is_locked": conv.is_locked, "locked_by": conv.locked_by, "is_direct": conv.is_direct, "is_archived": conv.is_archived, "is_pinned": is_pinned_by_user, "created_by": str(conv.created_by) if conv.created_by else None, "created_by_type": conv.created_by_type, "last_msg_at": conv.last_msg_at.isoformat() if conv.last_msg_at else None, "last_msg_preview": conv.last_msg_preview, "last_msg_sender_type": conv.last_msg_sender_type, "participants": [participant_to_response(p) for p in participants], "unread_count": unread_count, "metadata": conv.metadata_ or {}, } def participant_to_response(p: CommParticipant) -> dict[str, Any]: """Convert a CommParticipant to a response dict.""" return { "id": str(p.id), "conversation_id": str(p.conversation_id), "participant_id": str(p.participant_id) if p.participant_id else None, "participant_type": p.participant_type, "display_name": p.display_name, "role": p.role, "joined_at": p.joined_at.isoformat() if p.joined_at else None, } def message_to_response( msg: CommMessage, blocks: list[CommMessageBlock] | None = None, attachments: list[CommMessageAttachment] | None = None, reactions: list[CommMessageReaction] | None = None, ) -> dict[str, Any]: """Convert a CommMessage to a response dict.""" return { "id": str(msg.id), "conversation_id": str(msg.conversation_id), "sender_id": str(msg.sender_id) if msg.sender_id else None, "sender_type": msg.sender_type, "content": msg.content, "content_format": msg.content_format, "metadata": msg.metadata_ or {}, "reply_to_id": str(msg.reply_to_id) if msg.reply_to_id else None, "is_pinned": msg.is_pinned, "created_at": msg.created_at.isoformat() if msg.created_at else None, "edited_at": msg.edited_at.isoformat() if msg.edited_at else None, "blocks": [ { "id": str(b.id), "block_type": b.block_type, "block_data": b.block_data, "sort_order": b.sort_order, } for b in (blocks or []) ], "attachments": [ { "id": str(a.id), "file_id": str(a.file_id) if a.file_id else None, "file_source": a.file_source, "file_name": a.file_name, "file_type": a.file_type, "file_size": a.file_size, "thumbnail_path": a.thumbnail_path, } for a in (attachments or []) ], "reactions": [ { "id": str(r.id), "message_id": str(r.message_id), "user_id": str(r.user_id), "emoji": r.emoji, } for r in (reactions or []) ], } # ─── Conversation CRUD ─── async def list_conversations( db: AsyncSession, tenant_id: uuid.UUID, user_id: uuid.UUID, include_archived: bool = False, ) -> list[dict[str, Any]]: """List all conversations for a user.""" # Get conversations where user is a participant result = await db.execute( select(CommConversation) .join(CommParticipant, CommParticipant.conversation_id == CommConversation.id) .where( CommParticipant.participant_id == user_id, CommParticipant.participant_type == "user", CommParticipant.left_at.is_(None), CommConversation.tenant_id == tenant_id, CommConversation.deleted_at.is_(None), ) .order_by(CommConversation.last_msg_at.desc().nullslast()) ) conversations = result.scalars().all() # Get user's pinned conversations pins_result = await db.execute( select(CommConversationPin).where( CommConversationPin.user_id == user_id, CommConversationPin.tenant_id == tenant_id, ) ) pinned_ids = {p.conversation_id for p in pins_result.scalars().all()} conv_list = [] for conv in conversations: if conv.is_archived and not include_archived: continue # Get participants parts_result = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conv.id, CommParticipant.left_at.is_(None), ) ) participants = list(parts_result.scalars().all()) # Get unread count unread = await _get_unread_count(db, tenant_id, conv.id, user_id) conv_list.append( conversation_to_response( conv, participants, unread_count=unread, is_pinned_by_user=conv.id in pinned_ids ) ) # Sort: pinned first, then by last_msg_at conv_list.sort(key=lambda c: (not c["is_pinned"], c["last_msg_at"] or ""), reverse=False) # Actually: pinned first (True > False in reverse), then newest first conv_list.sort(key=lambda c: c["last_msg_at"] or "0000", reverse=True) conv_list.sort(key=lambda c: c["is_pinned"], reverse=True) return conv_list async def get_conversation( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, user_id: uuid.UUID, ) -> dict[str, Any] | None: """Get a single conversation with participants.""" result = await db.execute( select(CommConversation).where( CommConversation.id == conversation_id, CommConversation.tenant_id == tenant_id, CommConversation.deleted_at.is_(None), ) ) conv = result.scalar_one_or_none() if conv is None: return None # Check user is participant if not await CommRBAC.is_participant(db, conversation_id, user_id): return None parts_result = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conv.id, CommParticipant.left_at.is_(None), ) ) participants = list(parts_result.scalars().all()) # Check pinned pin_result = await db.execute( select(CommConversationPin).where( CommConversationPin.conversation_id == conv.id, CommConversationPin.user_id == user_id, ) ) is_pinned = pin_result.scalar_one_or_none() is not None unread = await _get_unread_count(db, tenant_id, conv.id, user_id) return conversation_to_response(conv, participants, unread_count=unread, is_pinned_by_user=is_pinned) async def create_conversation( db: AsyncSession, tenant_id: uuid.UUID, user_id: uuid.UUID, title: str | None = None, participant_ids: list[str] | None = None, is_direct: bool = False, initial_message: str | None = None, ) -> dict[str, Any]: """Create a new conversation.""" conv = CommConversation( tenant_id=tenant_id, title=title, is_direct=is_direct, created_by=user_id, created_by_type="user", metadata_={}, ) db.add(conv) await db.flush() # Add creator as admin creator = CommParticipant( tenant_id=tenant_id, conversation_id=conv.id, participant_id=user_id, participant_type="user", role="admin", ) db.add(creator) # Add other participants for pid_str in (participant_ids or []): try: pid = uuid.UUID(pid_str) if pid == user_id: continue p = CommParticipant( tenant_id=tenant_id, conversation_id=conv.id, participant_id=pid, participant_type="user", role="member", ) db.add(p) except ValueError: logger.warning(f"Invalid participant UUID: {pid_str}") await db.flush() # Send initial message if provided if initial_message: await send_message( db, tenant_id, conv.id, user_id, "user", content=initial_message, content_format="text", ) # Publish event event_bus = get_event_bus() await event_bus.publish("conversation.created", { "conversation_id": str(conv.id), "tenant_id": str(tenant_id), "created_by": str(user_id), }) # Get all participants for response parts_result = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conv.id, CommParticipant.left_at.is_(None), ) ) participants = list(parts_result.scalars().all()) return conversation_to_response(conv, participants) async def update_conversation( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, user_id: uuid.UUID, title: str | None = None, is_archived: bool | None = None, ) -> dict[str, Any] | None: """Update a conversation.""" result = await db.execute( select(CommConversation).where( CommConversation.id == conversation_id, CommConversation.tenant_id == tenant_id, CommConversation.deleted_at.is_(None), ) ) conv = result.scalar_one_or_none() if conv is None: return None # Check locked if conv.is_locked and title is not None: # Only the locking plugin can change title on locked conversations # Users cannot pass elif title is not None: conv.title = title conv.title_set_by = user_id if is_archived is not None: conv.is_archived = is_archived await db.flush() return await get_conversation(db, tenant_id, conversation_id, user_id) async def pin_conversation( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, user_id: uuid.UUID, ) -> bool: """Pin a conversation for a user.""" existing = await db.execute( select(CommConversationPin).where( CommConversationPin.conversation_id == conversation_id, CommConversationPin.user_id == user_id, ) ) if existing.scalar_one_or_none() is None: pin = CommConversationPin( tenant_id=tenant_id, conversation_id=conversation_id, user_id=user_id, ) db.add(pin) await db.flush() return True async def unpin_conversation( db: AsyncSession, conversation_id: uuid.UUID, user_id: uuid.UUID, ) -> bool: """Unpin a conversation for a user.""" result = await db.execute( select(CommConversationPin).where( CommConversationPin.conversation_id == conversation_id, CommConversationPin.user_id == user_id, ) ) pin = result.scalar_one_or_none() if pin: await db.delete(pin) await db.flush() return True async def mute_conversation( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, user_id: uuid.UUID, ) -> bool: """Mute a conversation for a user.""" existing = await db.execute( select(CommConversationMute).where( CommConversationMute.conversation_id == conversation_id, CommConversationMute.user_id == user_id, ) ) if existing.scalar_one_or_none() is None: mute = CommConversationMute( tenant_id=tenant_id, conversation_id=conversation_id, user_id=user_id, ) db.add(mute) await db.flush() return True async def unmute_conversation( db: AsyncSession, conversation_id: uuid.UUID, user_id: uuid.UUID, ) -> bool: """Unmute a conversation for a user.""" result = await db.execute( select(CommConversationMute).where( CommConversationMute.conversation_id == conversation_id, CommConversationMute.user_id == user_id, ) ) mute = result.scalar_one_or_none() if mute: await db.delete(mute) await db.flush() return True # ─── Participant Management ─── async def add_participant( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, participant_id: str, participant_type: str = "user", role: str = "member", display_name: str | None = None, ) -> dict[str, Any] | None: """Add a participant to a conversation.""" try: pid = uuid.UUID(participant_id) if participant_type == "user" else None except ValueError: return None existing = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conversation_id, CommParticipant.participant_id == pid if pid else CommParticipant.participant_type == participant_type, CommParticipant.participant_type == participant_type, CommParticipant.left_at.is_(None), ) ) if existing.scalar_one_or_none() is not None: return None # Already a participant p = CommParticipant( tenant_id=tenant_id, conversation_id=conversation_id, participant_id=pid, participant_type=participant_type, role=role, display_name=display_name, ) db.add(p) await db.flush() # Publish event event_bus = get_event_bus() await event_bus.publish("participant.joined", { "conversation_id": str(conversation_id), "participant_id": participant_id, "participant_type": participant_type, "tenant_id": str(tenant_id), }) return participant_to_response(p) async def remove_participant( db: AsyncSession, conversation_id: uuid.UUID, participant_id: uuid.UUID, ) -> bool: """Remove a participant from a conversation (set left_at).""" result = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conversation_id, CommParticipant.participant_id == participant_id, CommParticipant.participant_type == "user", CommParticipant.left_at.is_(None), ) ) p = result.scalar_one_or_none() if p is None: return False p.left_at = datetime.now(timezone.utc) await db.flush() event_bus = get_event_bus() await event_bus.publish("participant.left", { "conversation_id": str(conversation_id), "participant_id": str(participant_id), }) return True async def change_role( db: AsyncSession, conversation_id: uuid.UUID, participant_id: uuid.UUID, new_role: str, ) -> dict[str, Any] | None: """Change a participant's role.""" result = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conversation_id, CommParticipant.participant_id == participant_id, CommParticipant.participant_type == "user", CommParticipant.left_at.is_(None), ) ) p = result.scalar_one_or_none() if p is None: return None p.role = new_role await db.flush() return participant_to_response(p) # ─── Messages ─── async def get_messages( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, page: int = 1, page_size: int = 50, before_id: uuid.UUID | None = None, ) -> dict[str, Any]: """Get paginated messages for a conversation.""" query = select(CommMessage).where( CommMessage.conversation_id == conversation_id, CommMessage.tenant_id == tenant_id, CommMessage.deleted_at.is_(None), ).order_by(CommMessage.created_at.desc()) if before_id: before_msg = await db.execute( select(CommMessage).where(CommMessage.id == before_id) ) before = before_msg.scalar_one_or_none() if before: query = query.where(CommMessage.created_at < before.created_at) query = query.offset((page - 1) * page_size).limit(page_size) result = await db.execute(query) messages = list(result.scalars().all()) # Get blocks, attachments, reactions for each message msg_ids = [m.id for m in messages] blocks_map: dict[uuid.UUID, list] = {} attachments_map: dict[uuid.UUID, list] = {} reactions_map: dict[uuid.UUID, list] = {} if msg_ids: blocks_result = await db.execute( select(CommMessageBlock).where( CommMessageBlock.message_id.in_(msg_ids), CommMessageBlock.deleted_at.is_(None), ).order_by(CommMessageBlock.sort_order) ) for b in blocks_result.scalars().all(): blocks_map.setdefault(b.message_id, []).append(b) atts_result = await db.execute( select(CommMessageAttachment).where( CommMessageAttachment.message_id.in_(msg_ids), CommMessageAttachment.deleted_at.is_(None), ) ) for a in atts_result.scalars().all(): attachments_map.setdefault(a.message_id, []).append(a) reactions_result = await db.execute( select(CommMessageReaction).where( CommMessageReaction.message_id.in_(msg_ids), ) ) for r in reactions_result.scalars().all(): reactions_map.setdefault(r.message_id, []).append(r) items = [] for msg in reversed(messages): # chronological order items.append( message_to_response( msg, blocks=blocks_map.get(msg.id, []), attachments=attachments_map.get(msg.id, []), reactions=reactions_map.get(msg.id, []), ) ) # Total count count_result = await db.execute( select(func.count()).select_from(CommMessage).where( CommMessage.conversation_id == conversation_id, CommMessage.tenant_id == tenant_id, CommMessage.deleted_at.is_(None), ) ) total = count_result.scalar() or 0 has_more = (page * page_size) < total return {"items": items, "total": total, "page": page, "has_more": has_more} async def send_message( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, sender_id: uuid.UUID | None, sender_type: str, content: str = "", content_format: str = "text", blocks: list[dict[str, Any]] | None = None, reply_to_id: str | None = None, attachments: list[dict[str, Any]] | None = None, metadata: dict[str, Any] | None = None, trigger_depth: int = 0, ) -> dict[str, Any]: """Send a message to a conversation and trigger participant handlers.""" # Create message msg = CommMessage( tenant_id=tenant_id, conversation_id=conversation_id, sender_id=sender_id, sender_type=sender_type, content=content, content_format=content_format, metadata_=metadata or {}, ) if reply_to_id: try: msg.reply_to_id = uuid.UUID(reply_to_id) except ValueError: pass db.add(msg) await db.flush() # Create blocks if blocks: for i, block in enumerate(blocks): b = CommMessageBlock( tenant_id=tenant_id, message_id=msg.id, block_type=block["block_type"], block_data=block["block_data"], sort_order=i, ) db.add(b) # Create attachments if attachments: for att in attachments: a = CommMessageAttachment( tenant_id=tenant_id, message_id=msg.id, file_id=uuid.UUID(att["file_id"]) if att.get("file_id") else None, file_source=att.get("file_source", "comm"), file_name=att.get("file_name", ""), file_type=att.get("file_type", "application/octet-stream"), file_size=att.get("file_size"), ) db.add(a) await db.flush() # Update conversation last_msg await db.execute( update(CommConversation) .where(CommConversation.id == conversation_id) .values( last_msg_at=datetime.now(timezone.utc), last_msg_preview=content[:200] if content else "", last_msg_sender_type=sender_type, ) ) # Publish event event_bus = get_event_bus() await event_bus.publish("message.received", { "conversation_id": str(conversation_id), "message_id": str(msg.id), "sender_type": sender_type, "tenant_id": str(tenant_id), "content": content, "trigger_depth": trigger_depth, }) # Trigger participant handlers (if not at max depth) if trigger_depth < MAX_TRIGGER_DEPTH: await _trigger_participants( db, tenant_id, conversation_id, msg, trigger_depth ) # Load blocks/attachments/reactions for response blocks_result = await db.execute( select(CommMessageBlock).where( CommMessageBlock.message_id == msg.id, CommMessageBlock.deleted_at.is_(None), ).order_by(CommMessageBlock.sort_order) ) msg_blocks = list(blocks_result.scalars().all()) atts_result = await db.execute( select(CommMessageAttachment).where( CommMessageAttachment.message_id == msg.id, CommMessageAttachment.deleted_at.is_(None), ) ) msg_atts = list(atts_result.scalars().all()) return message_to_response(msg, blocks=msg_blocks, attachments=msg_atts) async def _trigger_participants( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, message: CommMessage, trigger_depth: int, ) -> None: """Trigger participant handlers for non-user participants.""" # Get conversation participants result = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conversation_id, CommParticipant.left_at.is_(None), CommParticipant.participant_type != "user", ) ) non_user_participants = list(result.scalars().all()) if not non_user_participants: return # Get conversation info conv_result = await db.execute( select(CommConversation).where(CommConversation.id == conversation_id) ) conv = conv_result.scalar_one_or_none() if conv is None: return # Parse mentions mentions = parse_mentions(message.content) # Build conversation dict all_parts_result = await db.execute( select(CommParticipant).where( CommParticipant.conversation_id == conversation_id, CommParticipant.left_at.is_(None), ) ) all_parts = [participant_to_response(p) for p in all_parts_result.scalars().all()] conv_dict = conversation_to_response(conv, []) conv_dict["participants"] = all_parts msg_dict = message_to_response(message) context = {"tenant_id": str(tenant_id), "trigger_depth": trigger_depth} registry = get_participant_registry() for p in non_user_participants: handler = registry.get_handler(p.participant_type) if handler is None: continue try: responses = await handler.on_message_received( conversation_id=conversation_id, message=msg_dict, conversation=conv_dict, mentions=mentions, context=context, ) if responses: for resp in responses: await send_message( db, tenant_id, conversation_id, sender_id=None, sender_type=p.participant_type, content=resp.get("content", ""), content_format=resp.get("content_format", "text"), blocks=resp.get("blocks"), metadata={ **(resp.get("metadata") or {}), "triggered_by": str(message.id), "trigger_depth": trigger_depth + 1, }, trigger_depth=trigger_depth + 1, ) except Exception: logger.exception( f"Participant handler error for type {p.participant_type}" ) async def edit_message( db: AsyncSession, tenant_id: uuid.UUID, message_id: uuid.UUID, user_id: uuid.UUID, new_content: str, ) -> dict[str, Any] | None: """Edit a message, storing the old version in history.""" result = await db.execute( select(CommMessage).where( CommMessage.id == message_id, CommMessage.tenant_id == tenant_id, CommMessage.deleted_at.is_(None), ) ) msg = result.scalar_one_or_none() if msg is None: return None # Get old blocks blocks_result = await db.execute( select(CommMessageBlock).where( CommMessageBlock.message_id == message_id, CommMessageBlock.deleted_at.is_(None), ) ) old_blocks = [b.block_data for b in blocks_result.scalars().all()] # Save edit history edit = CommMessageEdit( tenant_id=tenant_id, message_id=message_id, old_content=msg.content, old_blocks=old_blocks, edited_by=user_id, ) db.add(edit) # Update message msg.content = new_content msg.edited_at = datetime.now(timezone.utc) await db.flush() return message_to_response(msg) async def delete_message( db: AsyncSession, message_id: uuid.UUID, ) -> bool: """Soft-delete a message.""" result = await db.execute( select(CommMessage).where(CommMessage.id == message_id) ) msg = result.scalar_one_or_none() if msg is None: return False msg.deleted_at = datetime.now(timezone.utc) await db.flush() return True # ─── Reactions ─── async def add_reaction( db: AsyncSession, tenant_id: uuid.UUID, message_id: uuid.UUID, user_id: uuid.UUID, emoji: str, ) -> dict[str, Any] | None: """Add an emoji reaction to a message.""" existing = await db.execute( select(CommMessageReaction).where( CommMessageReaction.message_id == message_id, CommMessageReaction.user_id == user_id, CommMessageReaction.emoji == emoji, ) ) if existing.scalar_one_or_none() is not None: return None # Already reacted r = CommMessageReaction( tenant_id=tenant_id, message_id=message_id, user_id=user_id, emoji=emoji, ) db.add(r) await db.flush() event_bus = get_event_bus() await event_bus.publish("reaction.added", { "message_id": str(message_id), "emoji": emoji, "user_id": str(user_id), }) return { "id": str(r.id), "message_id": str(r.message_id), "user_id": str(r.user_id), "emoji": r.emoji, } async def remove_reaction( db: AsyncSession, message_id: uuid.UUID, user_id: uuid.UUID, emoji: str, ) -> bool: """Remove an emoji reaction.""" result = await db.execute( select(CommMessageReaction).where( CommMessageReaction.message_id == message_id, CommMessageReaction.user_id == user_id, CommMessageReaction.emoji == emoji, ) ) r = result.scalar_one_or_none() if r is None: return False await db.delete(r) await db.flush() return True # ─── Read State ─── async def mark_read( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, user_id: uuid.UUID, last_read_msg_id: str | None = None, ) -> bool: """Mark conversation as read up to a message.""" result = await db.execute( select(CommMessageRead).where( CommMessageRead.conversation_id == conversation_id, CommMessageRead.user_id == user_id, ) ) read = result.scalar_one_or_none() msg_id = uuid.UUID(last_read_msg_id) if last_read_msg_id else None if read is None: read = CommMessageRead( tenant_id=tenant_id, conversation_id=conversation_id, user_id=user_id, last_read_msg_id=msg_id, ) db.add(read) else: read.last_read_msg_id = msg_id read.last_read_at = datetime.now(timezone.utc) await db.flush() return True async def _get_unread_count( db: AsyncSession, tenant_id: uuid.UUID, conversation_id: uuid.UUID, user_id: uuid.UUID, ) -> int: """Get unread message count for a user in a conversation.""" # Get last read message read_result = await db.execute( select(CommMessageRead).where( CommMessageRead.conversation_id == conversation_id, CommMessageRead.user_id == user_id, ) ) read = read_result.scalar_one_or_none() query = select(func.count()).select_from(CommMessage).where( CommMessage.conversation_id == conversation_id, CommMessage.tenant_id == tenant_id, CommMessage.deleted_at.is_(None), CommMessage.sender_type != "system", # Don't count system messages? Or count all? ) if read and read.last_read_at: query = query.where(CommMessage.created_at > read.last_read_at) result = await db.execute(query) return result.scalar() or 0 # ─── Plugin Room Creation ─── async def create_plugin_room( db: AsyncSession, tenant_id: uuid.UUID, user_id: uuid.UUID, plugin_name: str, title: str, participant_type: str, user_role: str = "member", ) -> dict[str, Any]: """Create a locked, pinned room for a plugin (System, Live KI, Assistent). The room is locked (users can't change title/participants) and pinned for the user. """ # Check if room already exists for this user + plugin result = await db.execute( select(CommConversation).where( CommConversation.tenant_id == tenant_id, CommConversation.title == title, CommConversation.is_locked == True, CommConversation.locked_by == plugin_name, CommConversation.deleted_at.is_(None), ).join(CommParticipant, CommParticipant.conversation_id == CommConversation.id).where( CommParticipant.participant_id == user_id, CommParticipant.participant_type == "user", CommParticipant.left_at.is_(None), ) ) existing = result.scalar_one_or_none() if existing: # Already exists — return it return await get_conversation(db, tenant_id, existing.id, user_id) or {} # Create conversation conv = CommConversation( tenant_id=tenant_id, title=title, is_locked=True, locked_by=plugin_name, is_direct=False, created_by=None, created_by_type="plugin", metadata_={"plugin": plugin_name}, ) db.add(conv) await db.flush() # Add plugin as participant plugin_p = CommParticipant( tenant_id=tenant_id, conversation_id=conv.id, participant_id=None, participant_type=participant_type, role="admin", display_name=title, ) db.add(plugin_p) # Add user as participant user_p = CommParticipant( tenant_id=tenant_id, conversation_id=conv.id, participant_id=user_id, participant_type="user", role=user_role, ) db.add(user_p) # Pin for user pin = CommConversationPin( tenant_id=tenant_id, conversation_id=conv.id, user_id=user_id, ) db.add(pin) await db.flush() # Publish event event_bus = get_event_bus() await event_bus.publish("conversation.created", { "conversation_id": str(conv.id), "tenant_id": str(tenant_id), "created_by_type": "plugin", "plugin_name": plugin_name, }) parts = [plugin_p, user_p] return conversation_to_response(conv, parts, is_pinned_by_user=True)