Files
leocrm/app/plugins/builtins/kommunikation/routes.py
T

581 lines
22 KiB
Python
Raw Normal View History

"""REST API routes for the kommunikation plugin."""
from __future__ import annotations
import json
import logging
import uuid
from fastapi import (
APIRouter,
Depends,
File,
HTTPException,
Query,
UploadFile,
WebSocket,
WebSocketDisconnect,
)
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.db import get_db
from app.deps import get_current_user, require_permission
from app.plugins.builtins.kommunikation.content_types import list_block_types
from app.plugins.builtins.kommunikation.dms_bridge import DmsBridge
from app.plugins.builtins.kommunikation.rbac import CommRBAC
from app.plugins.builtins.kommunikation.schemas import (
ConversationCreate,
ConversationUpdate,
MessageCreate,
MessageUpdate,
MiniAppStartRequest,
ParticipantAdd,
ParticipantRoleUpdate,
ReactionCreate,
ReadStateUpdate,
)
from app.plugins.builtins.kommunikation.services import (
add_participant,
add_reaction,
change_role,
create_conversation,
delete_message,
edit_message,
get_conversation,
get_messages,
list_conversations,
mark_read,
mute_conversation,
pin_conversation,
remove_participant,
remove_reaction,
send_message,
unmute_conversation,
unpin_conversation,
update_conversation,
)
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/v1/comm", tags=["kommunikation"])
def _parse_uuid(val: str, field: str = "id") -> uuid.UUID:
try:
return uuid.UUID(val)
except (ValueError, TypeError):
raise HTTPException(400, detail={"detail": f"Invalid {field}", "code": "invalid_id"}) from None
# ─── Conversations ───
@router.get("/conversations")
async def list_user_conversations(
archived: bool = Query(False, description="Include archived conversations"),
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""List all conversations for the current user."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
convs = await list_conversations(db, tenant_id, user_id, include_archived=archived)
return {"items": convs, "total": len(convs)}
@router.post("/conversations")
async def create_new_conversation(
body: ConversationCreate,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Create a new conversation."""
if not await CommRBAC.can_user_create(current_user):
raise HTTPException(403, detail={"detail": "Permission denied", "code": "forbidden"})
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
return await create_conversation(
db, tenant_id, user_id,
title=body.title,
participant_ids=body.participant_ids,
is_direct=body.is_direct,
initial_message=body.initial_message,
2026-08-21 13:17:35 +02:00
metadata=body.metadata,
)
@router.get("/conversations/{conversation_id}")
async def get_single_conversation(
conversation_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Get a single conversation with participants."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
conv = await get_conversation(db, tenant_id, conv_id, user_id)
if conv is None:
raise HTTPException(404, detail={"detail": "Conversation not found", "code": "not_found"})
return conv
@router.patch("/conversations/{conversation_id}")
async def update_single_conversation(
conversation_id: str,
body: ConversationUpdate,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Update a conversation (title, archive)."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
conv = await update_conversation(db, tenant_id, conv_id, user_id, title=body.title, is_archived=body.is_archived)
if conv is None:
raise HTTPException(404, detail={"detail": "Conversation not found", "code": "not_found"})
return conv
@router.delete("/conversations/{conversation_id}")
async def leave_or_delete_conversation(
conversation_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Leave (member) or delete (admin) a conversation."""
user_id = uuid.UUID(current_user["user_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
# For now: just leave (set left_at)
success = await remove_participant(db, conv_id, user_id)
if not success:
raise HTTPException(404, detail={"detail": "Not a participant", "code": "not_found"})
return {"success": True}
@router.post("/conversations/{conversation_id}/pin")
async def pin_conv(
conversation_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Pin a conversation for the current user."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
await pin_conversation(db, tenant_id, conv_id, user_id)
return {"success": True}
@router.delete("/conversations/{conversation_id}/pin")
async def unpin_conv(
conversation_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Unpin a conversation."""
user_id = uuid.UUID(current_user["user_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
await unpin_conversation(db, conv_id, user_id)
return {"success": True}
@router.post("/conversations/{conversation_id}/mute")
async def mute_conv(
conversation_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Mute a conversation."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
await mute_conversation(db, tenant_id, conv_id, user_id)
return {"success": True}
@router.delete("/conversations/{conversation_id}/mute")
async def unmute_conv(
conversation_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Unmute a conversation."""
user_id = uuid.UUID(current_user["user_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
await unmute_conversation(db, conv_id, user_id)
return {"success": True}
# ─── Participants ───
@router.post("/conversations/{conversation_id}/participants", dependencies=[Depends(require_permission("comm:manage"))])
async def add_participant_endpoint(
conversation_id: str,
body: ParticipantAdd,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Add a participant to a conversation."""
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
if not await CommRBAC.can_user_manage(db, conv_id, current_user):
raise HTTPException(403, detail={"detail": "Cannot manage participants", "code": "forbidden"})
result = await add_participant(db, tenant_id, conv_id, body.participant_id, body.participant_type, body.role)
if result is None:
raise HTTPException(400, detail={"detail": "Already a participant or invalid ID", "code": "bad_request"})
return result
@router.delete("/conversations/{conversation_id}/participants/{participant_id}", dependencies=[Depends(require_permission("comm:manage"))])
async def remove_participant_endpoint(
conversation_id: str,
participant_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Remove a participant from a conversation."""
conv_id = _parse_uuid(conversation_id, "conversation_id")
pid = _parse_uuid(participant_id, "participant_id")
# Self-leave is always allowed
if pid != uuid.UUID(current_user["user_id"]):
if not await CommRBAC.can_user_manage(db, conv_id, current_user):
raise HTTPException(403, detail={"detail": "Cannot manage participants", "code": "forbidden"})
success = await remove_participant(db, conv_id, pid)
if not success:
raise HTTPException(404, detail={"detail": "Participant not found", "code": "not_found"})
return {"success": True}
@router.patch("/conversations/{conversation_id}/participants/{participant_id}", dependencies=[Depends(require_permission("comm:manage"))])
async def change_participant_role(
conversation_id: str,
participant_id: str,
body: ParticipantRoleUpdate,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Change a participant's role."""
conv_id = _parse_uuid(conversation_id, "conversation_id")
pid = _parse_uuid(participant_id, "participant_id")
if not await CommRBAC.can_user_change_role(db, conv_id, current_user):
raise HTTPException(403, detail={"detail": "Cannot change roles", "code": "forbidden"})
result = await change_role(db, conv_id, pid, body.role)
if result is None:
raise HTTPException(404, detail={"detail": "Participant not found", "code": "not_found"})
return result
# ─── Messages ───
@router.get("/conversations/{conversation_id}/messages", dependencies=[Depends(require_permission("comm:read"))])
async def get_conv_messages(
conversation_id: str,
page: int = Query(1, ge=1),
page_size: int = Query(50, ge=1, le=100),
before: str | None = Query(None),
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Get paginated messages for a conversation."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
if not await CommRBAC.is_participant(db, conv_id, user_id):
raise HTTPException(403, detail={"detail": "Not a participant", "code": "forbidden"})
before_id = _parse_uuid(before, "before") if before else None
return await get_messages(db, tenant_id, conv_id, page, page_size, before_id)
@router.post("/conversations/{conversation_id}/messages", dependencies=[Depends(require_permission("comm:write"))])
async def send_conv_message(
conversation_id: str,
body: MessageCreate,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Send a message to a conversation."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
if not await CommRBAC.can_user_write(db, conv_id, current_user):
raise HTTPException(403, detail={"detail": "Cannot write to this conversation", "code": "forbidden"})
blocks_data = [b.model_dump() for b in body.blocks] if body.blocks else None
attachments_data = [a.model_dump() for a in body.attachments] if body.attachments else None
return await send_message(
db, tenant_id, conv_id, user_id, "user",
content=body.content,
content_format=body.content_format,
blocks=blocks_data,
reply_to_id=body.reply_to_id,
attachments=attachments_data,
)
@router.patch("/messages/{message_id}", dependencies=[Depends(require_permission("comm:write"))])
async def update_msg(
message_id: str,
body: MessageUpdate,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Edit a message or mark as read."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
msg_id = _parse_uuid(message_id, "message_id")
if body.content is not None:
result = await edit_message(db, tenant_id, msg_id, user_id, body.content)
if result is None:
raise HTTPException(404, detail={"detail": "Message not found", "code": "not_found"})
return result
return {"success": True}
@router.delete("/messages/{message_id}", dependencies=[Depends(require_permission("comm:delete"))])
async def delete_msg(
message_id: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Delete a message."""
msg_id = _parse_uuid(message_id, "message_id")
success = await delete_message(db, msg_id)
if not success:
raise HTTPException(404, detail={"detail": "Message not found", "code": "not_found"})
return {"success": True}
# ─── Attachments ───
@router.post("/messages/{message_id}/attachments", response_model=None, dependencies=[Depends(require_permission("comm:write"))])
async def upload_attachment(
message_id: str,
file: UploadFile = File(...),
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Upload a file attachment to a message."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
msg_id = _parse_uuid(message_id, "message_id")
# Rate limit — UPLOAD policy
from app.core.rate_limit import RateLimitPolicy, check_rate_limit_policy
await check_rate_limit_policy(
f"rate:upload:comm:{tenant_id}:{user_id}",
RateLimitPolicy.UPLOAD,
)
# Get conversation_id from message
from sqlalchemy import select
from app.plugins.builtins.kommunikation.models import CommMessage
result = await db.execute(select(CommMessage).where(CommMessage.id == msg_id))
msg = result.scalar_one_or_none()
if msg is None:
raise HTTPException(404, detail={"detail": "Message not found", "code": "not_found"})
try:
return await DmsBridge.store_attachment(db, tenant_id, msg.conversation_id, user_id, file)
except ValueError as e:
raise HTTPException(413, detail={"detail": str(e), "code": "file_too_large"}) from e
# ─── Reactions ───
@router.post("/messages/{message_id}/reactions", dependencies=[Depends(require_permission("comm:write"))])
async def add_msg_reaction(
message_id: str,
body: ReactionCreate,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Add an emoji reaction to a message."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
msg_id = _parse_uuid(message_id, "message_id")
result = await add_reaction(db, tenant_id, msg_id, user_id, body.emoji)
if result is None:
raise HTTPException(409, detail={"detail": "Already reacted", "code": "conflict"})
return result
@router.delete("/messages/{message_id}/reactions/{emoji}", dependencies=[Depends(require_permission("comm:write"))])
async def remove_msg_reaction(
message_id: str,
emoji: str,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Remove an emoji reaction."""
user_id = uuid.UUID(current_user["user_id"])
msg_id = _parse_uuid(message_id, "message_id")
success = await remove_reaction(db, msg_id, user_id, emoji)
if not success:
raise HTTPException(404, detail={"detail": "Reaction not found", "code": "not_found"})
return {"success": True}
# ─── Read State ───
@router.post("/conversations/{conversation_id}/read", dependencies=[Depends(require_permission("comm:write"))])
async def mark_conv_read(
conversation_id: str,
body: ReadStateUpdate,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Mark a conversation as read."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
await mark_read(db, tenant_id, conv_id, user_id, body.last_read_msg_id)
return {"success": True}
# ─── Mini-Apps ───
@router.get("/miniapps", dependencies=[Depends(require_permission("comm:read"))])
async def list_miniapps(
current_user: dict = Depends(get_current_user),
):
"""List available mini-apps."""
from app.core.service_container import get_container
container = get_container()
if not container.has("comm_miniapps"):
return {"items": []}
registry = container.get("comm_miniapps")
return {"items": registry.list_apps()}
@router.post("/conversations/{conversation_id}/miniapps", dependencies=[Depends(require_permission("comm:write"))])
async def start_miniapp(
conversation_id: str,
body: MiniAppStartRequest,
current_user: dict = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Start a mini-app in a conversation."""
user_id = uuid.UUID(current_user["user_id"])
tenant_id = uuid.UUID(current_user["tenant_id"])
conv_id = _parse_uuid(conversation_id, "conversation_id")
if not await CommRBAC.can_user_write(db, conv_id, current_user):
raise HTTPException(403, detail={"detail": "Cannot write", "code": "forbidden"})
# Create a miniapp block as a message
return await send_message(
db, tenant_id, conv_id, user_id, "user",
content=f"[Mini-App: {body.app_id}]",
blocks=[{
"block_type": "miniapp",
"block_data": {"app_id": body.app_id, "config": body.config},
}],
)
# ─── Content Types ───
@router.get("/block-types", dependencies=[Depends(require_permission("comm:read"))])
async def get_block_types():
"""List all known content block types."""
return {"items": list_block_types()}
# ─── WebSocket ───
@router.websocket("/ws")
async def websocket_endpoint(
websocket: WebSocket,
):
"""WebSocket endpoint for real-time messaging.
Authenticates via session cookie. On connect, subscribes user to all their conversations.
"""
from app.core.db import async_session_maker
from app.core.service_container import get_container
# Get WebSocket manager from service container
container = get_container()
if not container.has("comm_websocket"):
await websocket.close(code=4003, reason="Messaging not available")
return
ws_manager = container.get("comm_websocket")
# connect() performs origin check, session auth, tenant check
async with async_session_maker() as db:
auth = await ws_manager.connect(websocket, db)
if auth is None:
return # connection was rejected and closed by ws_helpers
user_id = auth["user_id"]
tenant_id = auth["tenant_id"]
# Plugin-Gate: check if kommunikation plugin is active (global + tenant)
import uuid as _uuid
from sqlalchemy import text as sa_text
from app.core.permission_registry import get_permission_registry
try:
registry = get_permission_registry()
if not registry.is_plugin_active("kommunikation"):
await websocket.close(code=4003, reason="Plugin not active")
return
async with async_session_maker() as db:
result = await db.execute(
sa_text("SELECT is_active FROM tenant_plugin_activation WHERE plugin_name = :name AND tenant_id = :tid"),
{"name": "kommunikation", "tid": _uuid.UUID(tenant_id)},
)
row = result.first()
if row is not None and not row[0]:
await websocket.close(code=4003, reason="Plugin not active for tenant")
return
except Exception:
await websocket.close(code=4003, reason="Plugin check failed")
return
try:
while True:
data = await websocket.receive_text()
msg = json.loads(data)
msg_type = msg.get("type")
if msg_type == "ping":
await ws_manager.send_to_user(user_id, {"type": "pong"})
elif msg_type == "subscribe":
conv_id = msg.get("conversation_id")
if conv_id:
# P1.9 fix: Check if user is a participant of this conversation
from sqlalchemy import text as sql_text
from app.core.db import async_session_maker
try:
async with async_session_maker() as db:
await db.execute(sql_text("SELECT set_config('app.current_tenant_id', :tid, true)"), {"tid": tenant_id})
result = await db.execute(
sql_text("SELECT 1 FROM conversation_participants WHERE conversation_id = :cid AND user_id = :uid"),
{"cid": conv_id, "uid": user_id},
)
if result.first():
ws_manager.subscribe(conv_id, user_id)
else:
await ws_manager.send_to_user(user_id, {"type": "error", "message": "Not a participant of this conversation"})
except Exception:
logger.warning("Failed to check conversation participation for %s in %s", user_id, conv_id)
await ws_manager.send_to_user(user_id, {"type": "error", "message": "Cannot verify participation"})
elif msg_type == "unsubscribe":
conv_id = msg.get("conversation_id")
if conv_id:
ws_manager.unsubscribe(conv_id, user_id)
elif msg_type == "typing":
conv_id = msg.get("conversation_id")
is_typing = msg.get("is_typing", False)
if conv_id:
await ws_manager.send_to_conversation(
conv_id,
{"type": "typing", "conversation_id": conv_id, "user_id": user_id, "is_typing": is_typing},
exclude_user=user_id,
)
except WebSocketDisconnect:
await ws_manager.disconnect(websocket, user_id)
except Exception:
logger.exception("WebSocket error")
await ws_manager.disconnect(websocket, user_id)