diff --git a/app/ai/agent_loop.py b/app/ai/agent_loop.py index e3854c7..f376e94 100644 --- a/app/ai/agent_loop.py +++ b/app/ai/agent_loop.py @@ -394,27 +394,21 @@ async def run_react_loop( if agent_run_id: try: from app.plugins.builtins.contracts import get_contract_registry - from app.plugins.builtins.kommunikation.models import CommConversation - from sqlalchemy import select as sa_select komm = get_contract_registry().get("kommunikation") if komm: agent_id = getattr(agent_definition, "id", uuid.uuid4()) room_title = f"Agent: {getattr(agent_definition, 'name', 'Agent')}" - existing = await db.execute( - sa_select(CommConversation).where( - CommConversation.tenant_id == tenant_id, - CommConversation.title == room_title, - CommConversation.is_locked.is_(True), - CommConversation.locked_by == "automation", - CommConversation.deleted_at.is_(None), - ) + conv_id = await komm.find_locked_room_id( + db=db, + tenant_id=tenant_id, + plugin_name="automation", + title=room_title, ) - conv = existing.scalar_one_or_none() - if conv: + if conv_id: await komm.send_message( db=db, tenant_id=tenant_id, - conversation_id=conv.id, + conversation_id=conv_id, sender_id=agent_id, sender_type="agent", content=f"Approval required for tool '{tool_name}'", diff --git a/app/plugins/builtins/kommunikation/contracts.py b/app/plugins/builtins/kommunikation/contracts.py index 1d9e2ec..b94cd08 100644 --- a/app/plugins/builtins/kommunikation/contracts.py +++ b/app/plugins/builtins/kommunikation/contracts.py @@ -31,6 +31,7 @@ from app.plugins.builtins.kommunikation.participant_registry import ( ) from app.plugins.builtins.kommunikation.services import ( create_plugin_room, + find_locked_room_id, get_conversation, get_messages, parse_mentions, @@ -55,6 +56,7 @@ class KommunikationContract: get_messages = staticmethod(get_messages) send_message = staticmethod(send_message) create_plugin_room = staticmethod(create_plugin_room) + find_locked_room_id = staticmethod(find_locked_room_id) # ─── participant registry ─── get_participant_registry = staticmethod(get_participant_registry) @@ -92,6 +94,7 @@ __all__ = [ "get_messages", "send_message", "create_plugin_room", + "find_locked_room_id", "CommConversation", "CommMessage", "CommParticipant", diff --git a/app/plugins/builtins/kommunikation/services.py b/app/plugins/builtins/kommunikation/services.py index c32fe0c..ffe0b93 100644 --- a/app/plugins/builtins/kommunikation/services.py +++ b/app/plugins/builtins/kommunikation/services.py @@ -1047,6 +1047,30 @@ async def _get_unread_count( # ─── Plugin Room Creation ─── +async def find_locked_room_id( + db: AsyncSession, + tenant_id: uuid.UUID, + plugin_name: str, + title: str, +) -> uuid.UUID | None: + """Find the conversation ID of a locked plugin room by tenant and title. + + Matches the same room semantics as ``create_plugin_room``: locked rooms + are owned by the plugin (``locked_by == plugin_name``) and soft-deleted + conversations are excluded. Returns ``None`` when no room exists. + """ + result = await db.execute( + select(CommConversation.id).where( + CommConversation.tenant_id == tenant_id, + CommConversation.title == title, + CommConversation.is_locked.is_(True), + CommConversation.locked_by == plugin_name, + CommConversation.deleted_at.is_(None), + ) + ) + return result.scalar_one_or_none() + + async def create_plugin_room( db: AsyncSession, tenant_id: uuid.UUID, diff --git a/app/workflows/engine.py b/app/workflows/engine.py index 29b5175..e9e088b 100644 --- a/app/workflows/engine.py +++ b/app/workflows/engine.py @@ -92,22 +92,16 @@ class WorkflowEngine: # Post workflow completion to Communication (G-WORK) try: from app.plugins.builtins.contracts import get_contract_registry - from app.plugins.builtins.kommunikation.models import CommConversation - from sqlalchemy import select as sa_select komm = get_contract_registry().get("kommunikation") if komm and instance.initiated_by: room_title = f"Workflow: {workflow.name if hasattr(workflow, 'name') else str(instance.workflow_id)}" - existing = await self.db.execute( - sa_select(CommConversation).where( - CommConversation.tenant_id == self.tenant_id, - CommConversation.title == room_title, - CommConversation.is_locked.is_(True), - CommConversation.locked_by == "workflow", - CommConversation.deleted_at.is_(None), - ) + conv_id = await komm.find_locked_room_id( + db=self.db, + tenant_id=self.tenant_id, + plugin_name="workflow", + title=room_title, ) - conv = existing.scalar_one_or_none() - if not conv: + if not conv_id: room = await komm.create_plugin_room( db=self.db, tenant_id=self.tenant_id, @@ -117,8 +111,6 @@ class WorkflowEngine: participant_type="workflow", ) conv_id = uuid.UUID(room["conversation_id"]) - else: - conv_id = conv.id await komm.send_message( db=self.db, tenant_id=self.tenant_id,