# Conflicts: # apps/chat/serializers/messages.py # apps/chat/services/message.py # apps/chat/tests/test_messages.py # apps/chat/views/messages.py
140 lines
5 KiB
Python
140 lines
5 KiB
Python
from datetime import datetime, timezone
|
|
from uuid import UUID
|
|
|
|
from apps.chat.events.event import MessageSentEvent
|
|
from apps.chat.events.publishers.longpoll import LongPollPublisher
|
|
from apps.chat.events.publishers.push import PushPublisher
|
|
from apps.chat.events.publishers.websocket import WebSocketPublisher
|
|
from apps.chat.exceptions import ConversationClosedError
|
|
from apps.chat.integrations.mattermost.client import MattermostClient
|
|
from apps.chat.models import Conversation, ConversationStatus, MattermostAccountMapping
|
|
from apps.chat.services.account import AccountService
|
|
from apps.chat.services.storage import StorageService
|
|
|
|
MEDIA_MESSAGE_TYPES = {"image", "video", "voice"}
|
|
|
|
|
|
class MessageService:
|
|
def __init__(
|
|
self,
|
|
account_service: AccountService | None = None,
|
|
storage_service: StorageService | None = None,
|
|
mattermost_client: MattermostClient | None = None,
|
|
publishers: list | None = None,
|
|
):
|
|
self._account = account_service or AccountService()
|
|
self._storage_override = storage_service # lazy: only instantiate when needed
|
|
self._mm = mattermost_client or MattermostClient()
|
|
self._publishers = (
|
|
publishers
|
|
if publishers is not None
|
|
else [WebSocketPublisher(), PushPublisher(), LongPollPublisher()]
|
|
)
|
|
|
|
@property
|
|
def _storage(self) -> StorageService:
|
|
if self._storage_override is None:
|
|
self._storage_override = StorageService()
|
|
return self._storage_override
|
|
|
|
def send(
|
|
self,
|
|
conversation_uuid,
|
|
sender_uuid: UUID,
|
|
message_type: str,
|
|
text: str | None = None,
|
|
object_key: str | None = None,
|
|
) -> dict:
|
|
conversation = Conversation.objects.get(uuid=conversation_uuid)
|
|
if conversation.status == ConversationStatus.CLOSED:
|
|
raise ConversationClosedError(f"Conversation {conversation_uuid} is closed")
|
|
|
|
if not self._account.validate_user(sender_uuid):
|
|
raise ValueError(f"Sender {sender_uuid} is not valid")
|
|
|
|
# Media was already uploaded straight to MinIO by the client; the
|
|
# Mattermost post body only ever carries "<type>:<object_key>", never
|
|
# a URL, so a stored link can't outlive its presigned expiry.
|
|
download_url: str | None = None
|
|
if message_type in MEDIA_MESSAGE_TYPES:
|
|
mm_message = f"{message_type}:{object_key}"
|
|
download_url = self._storage.get_download_url(object_key)
|
|
else:
|
|
mm_message = text or ""
|
|
|
|
post_id = self._mm.post_message(conversation.mattermost_channel_id, mm_message)
|
|
|
|
event = MessageSentEvent(
|
|
chat_uuid=conversation_uuid,
|
|
post_id=post_id,
|
|
sender_uuid=sender_uuid,
|
|
message_type=message_type,
|
|
payload={"text": text, "object_key": object_key, "url": download_url},
|
|
)
|
|
for publisher in self._publishers:
|
|
publisher.publish(event)
|
|
|
|
return {
|
|
"post_id": post_id,
|
|
"sender_uuid": sender_uuid,
|
|
"message_type": message_type,
|
|
"text": text,
|
|
"url": download_url,
|
|
"created_at": None,
|
|
}
|
|
|
|
def list_messages(
|
|
self,
|
|
conversation_uuid,
|
|
page: int,
|
|
per_page: int,
|
|
since: int | None = None,
|
|
) -> list[dict]:
|
|
conversation = Conversation.objects.get(uuid=conversation_uuid)
|
|
raw_posts = self._mm.get_posts(
|
|
conversation.mattermost_channel_id,
|
|
page=page,
|
|
per_page=per_page,
|
|
since=since,
|
|
)
|
|
|
|
mm_user_ids = {p["user_id"] for p in raw_posts if p.get("user_id")}
|
|
mappings = {
|
|
m.mattermost_user_id: m.user_uuid
|
|
for m in MattermostAccountMapping.objects.filter(
|
|
mattermost_user_id__in=mm_user_ids
|
|
)
|
|
} if mm_user_ids else {}
|
|
|
|
return [self._normalize_post(p, mappings) for p in raw_posts]
|
|
|
|
def _normalize_post(self, post: dict, mappings: dict) -> dict:
|
|
mm_uid = post.get("user_id")
|
|
sender_uuid = mappings.get(mm_uid)
|
|
msg = post.get("message", "")
|
|
created_ms = post.get("create_at")
|
|
created_at = (
|
|
datetime.fromtimestamp(created_ms / 1000, tz=timezone.utc)
|
|
if created_ms
|
|
else None
|
|
)
|
|
|
|
message_type, sep, object_key = msg.partition(":")
|
|
if sep and message_type in MEDIA_MESSAGE_TYPES and object_key:
|
|
return {
|
|
"post_id": post.get("id"),
|
|
"sender_uuid": sender_uuid,
|
|
"message_type": message_type,
|
|
"text": None,
|
|
"url": self._storage.get_download_url(object_key),
|
|
"created_at": created_at,
|
|
}
|
|
|
|
return {
|
|
"post_id": post.get("id"),
|
|
"sender_uuid": sender_uuid,
|
|
"message_type": "text",
|
|
"text": msg,
|
|
"url": None,
|
|
"created_at": created_at,
|
|
}
|