diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..18b410d --- /dev/null +++ b/.env.example @@ -0,0 +1,41 @@ +# Database +DB_NAME=gooyal_chat +DB_USER=postgres +DB_HOST=127.0.0.1 +DB_PASSWORD= +DB_PORT=5432 + +DEBUG=true + +# OAuth2 / SSO provider +OAUTH2_PROVIDER_BASE_PUBLIC_URL=https://accounts.example.com/oauth2 +OAUTH2_PROVIDER_BASE_PRIVATE_URL=https://accounts.example.com/oauth2 +OAUTH2_PROVIDER_CLIENT_ID= +OAUTH2_PROVIDER_CLIENT_SECRET= +OAUTH2_PROVIDER_SCOPES= + +# Redis (used for cache, channel layer, and long-poll) +REDIS_BASE_URL=redis://127.0.0.1:6379/1 + +# Mattermost +# MATTERMOST_TOKEN must be a Personal Access Token (not a session token — session tokens expire) +# Enable at: System Console → Integrations → Integration Management → Enable Personal Access Tokens +MATTERMOST_URL=https://chat.example.com +MATTERMOST_TOKEN= +MATTERMOST_TEAM_ID= +MATTERMOST_SERVICE_USERNAME=chat-service +MATTERMOST_SERVICE_EMAIL=chat-service@local.invalid + +# MinIO object storage +MINIO_ENDPOINT=localhost:9000 +MINIO_ACCESS_KEY=minioadmin +MINIO_SECRET_KEY=minioadmin +MINIO_BUCKET_CHAT=chat + +# Chat long-poll tuning +CHAT_LONG_POLL_TIMEOUT_SECONDS=25 +CHAT_LONG_POLL_INTERVAL_SECONDS=1 + +# Native library paths (macOS Homebrew — leave empty on Linux) +GDAL_LIBRARY_PATH= +GEOS_LIBRARY_PATH= diff --git a/.gitignore b/.gitignore index c25a548..926e870 100644 --- a/.gitignore +++ b/.gitignore @@ -5,3 +5,8 @@ media /delme.py /log/accounts.log /log/errors.log +**/__pycache__ +**/*.pyc +**/*.pyo +.claude +/log \ No newline at end of file diff --git a/apps/chat/apps.py b/apps/chat/apps.py index 45f098c..401333b 100644 --- a/apps/chat/apps.py +++ b/apps/chat/apps.py @@ -3,3 +3,4 @@ from django.apps import AppConfig class ChatConfig(AppConfig): name = 'apps.chat' + default_auto_field = 'django.db.models.BigAutoField' diff --git a/apps/chat/events/__init__.py b/apps/chat/events/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/events/base.py b/apps/chat/events/base.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/events/event.py b/apps/chat/events/event.py new file mode 100644 index 0000000..85871c5 --- /dev/null +++ b/apps/chat/events/event.py @@ -0,0 +1,11 @@ +from dataclasses import dataclass +from uuid import UUID + + +@dataclass +class MessageSentEvent: + chat_uuid: UUID + post_id: str + sender_uuid: UUID + message_type: str + payload: dict diff --git a/apps/chat/events/publishers/__init__.py b/apps/chat/events/publishers/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/events/publishers/longpoll.py b/apps/chat/events/publishers/longpoll.py new file mode 100644 index 0000000..0afd568 --- /dev/null +++ b/apps/chat/events/publishers/longpoll.py @@ -0,0 +1,23 @@ +import json + +from django_redis import get_redis_connection + +from apps.chat.events.event import MessageSentEvent + +_KEY_TTL = 300 + + +class LongPollPublisher: + def publish(self, event: MessageSentEvent) -> None: + redis = get_redis_connection("default") + key = f"chat_events:{event.chat_uuid}" + payload = json.dumps( + { + "post_id": event.post_id, + "sender_uuid": str(event.sender_uuid), + "message_type": event.message_type, + **{k: v for k, v in event.payload.items() if v is not None}, + } + ) + redis.rpush(key, payload) + redis.expire(key, _KEY_TTL) diff --git a/apps/chat/events/publishers/push.py b/apps/chat/events/publishers/push.py new file mode 100644 index 0000000..fbba604 --- /dev/null +++ b/apps/chat/events/publishers/push.py @@ -0,0 +1,7 @@ +from apps.chat.events.event import MessageSentEvent + + +class PushPublisher: + # Future: FCM / APNS integration + def publish(self, event: MessageSentEvent) -> None: + pass diff --git a/apps/chat/events/publishers/websocket.py b/apps/chat/events/publishers/websocket.py new file mode 100644 index 0000000..934d8f1 --- /dev/null +++ b/apps/chat/events/publishers/websocket.py @@ -0,0 +1,20 @@ +from asgiref.sync import async_to_sync +from channels.layers import get_channel_layer + +from apps.chat.events.event import MessageSentEvent + + +class WebSocketPublisher: + def publish(self, event: MessageSentEvent) -> None: + channel_layer = get_channel_layer() + group_name = f"chat_{event.chat_uuid}" + async_to_sync(channel_layer.group_send)( + group_name, + { + "type": "chat.message", + "post_id": event.post_id, + "sender_uuid": str(event.sender_uuid), + "message_type": event.message_type, + **event.payload, + }, + ) diff --git a/apps/chat/exceptions.py b/apps/chat/exceptions.py new file mode 100644 index 0000000..53b452c --- /dev/null +++ b/apps/chat/exceptions.py @@ -0,0 +1,2 @@ +class ConversationClosedError(Exception): + """Raised when attempting to send a message to a closed conversation.""" diff --git a/apps/chat/integrations/__init__.py b/apps/chat/integrations/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/integrations/mattermost/__init__.py b/apps/chat/integrations/mattermost/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/integrations/mattermost/client.py b/apps/chat/integrations/mattermost/client.py new file mode 100644 index 0000000..659a5ed --- /dev/null +++ b/apps/chat/integrations/mattermost/client.py @@ -0,0 +1,218 @@ +import secrets +import string +from urllib.parse import urlparse +from uuid import UUID + +from django.conf import settings +from mattermostdriver import Driver +from mattermostdriver.exceptions import ResourceNotFound + +from apps.chat.integrations.mattermost.exceptions import MattermostError + + +class MattermostClient: + """ + Thin adapter around the Mattermost API. + + All calls use a single service-level admin token — no per-user tokens. + Mattermost internal IDs (user IDs, channel IDs, post IDs) are opaque + strings from the perspective of callers; they must never appear in public + API responses or service-layer interfaces. + """ + + def __init__(self): + url = getattr(settings, "MATTERMOST_URL", None) + token = getattr(settings, "MATTERMOST_TOKEN", None) + + if not url: + raise MattermostError("MATTERMOST_URL is not configured") + if not token: + raise MattermostError("MATTERMOST_TOKEN is not configured") + + parsed = urlparse(url) + scheme = parsed.scheme or "https" + hostname = parsed.hostname or url + port = parsed.port or (443 if scheme == "https" else 80) + + self._team_id: str | None = getattr(settings, "MATTERMOST_TEAM_ID", None) or None + + self._driver = Driver( + { + "url": hostname, + "token": token, + "scheme": scheme, + "port": port, + "verify": True, + "debug": False, + } + ) + self._driver.login() + + # ------------------------------------------------------------------ + # Internal helpers — username and email are deterministic from UUID + # so the same user always maps to the same Mattermost account. + # ------------------------------------------------------------------ + + def _username_for(self, user_id: UUID) -> str: + return f"chat-{UUID(str(user_id)).hex}" + + def _email_for(self, user_id: UUID) -> str: + return f"chat-{UUID(str(user_id)).hex}@mm.internal" + + @staticmethod + def _random_password() -> str: + alphabet = string.ascii_letters + string.digits + "!@#$%" + return "".join(secrets.choice(alphabet) for _ in range(32)) + + # ------------------------------------------------------------------ + # Public API — return only opaque strings, never raw Mattermost dicts + # with identifying fields surfaced beyond this layer. + # ------------------------------------------------------------------ + + def _ensure_team_member(self, mm_user_id: str) -> None: + if not self._team_id: + return + try: + self._driver.teams.add_user_to_team( + self._team_id, options={"team_id": self._team_id, "user_id": mm_user_id} + ) + except Exception: + pass # already a member or team not configured — not fatal + + def get_or_create_user(self, user_id: UUID) -> str: + """ + Ensure a Mattermost account exists for the given external UUID. + Returns the opaque mattermost_user_id string. + """ + username = self._username_for(user_id) + + try: + user = self._driver.users.get_user_by_username(username) + mm_user_id = user["id"] + self._ensure_team_member(mm_user_id) + return mm_user_id + except ResourceNotFound: + pass + except Exception as exc: + raise MattermostError(f"User lookup failed: {exc}") from exc + + try: + user = self._driver.users.create_user( + options={ + "username": username, + "email": self._email_for(user_id), + "password": self._random_password(), + } + ) + mm_user_id = user["id"] + self._ensure_team_member(mm_user_id) + return mm_user_id + except Exception as exc: + # Guard against a concurrent request having created the account + # between our lookup and our create attempt. + try: + user = self._driver.users.get_user_by_username(username) + mm_user_id = user["id"] + self._ensure_team_member(mm_user_id) + return mm_user_id + except Exception: + pass + raise MattermostError(f"User creation failed: {exc}") from exc + + def create_private_channel(self, mm_user_ids: list[str]) -> str: + """ + Create a new private Mattermost channel and add all members. + Each call always creates a brand-new channel. + Returns the opaque mattermost_channel_id string. + """ + if not self._team_id: + raise MattermostError("MATTERMOST_TEAM_ID is not configured") + + channel_name = f"conv-{secrets.token_hex(8)}" + + try: + channel = self._driver.channels.create_channel( + options={ + "team_id": self._team_id, + "name": channel_name, + "display_name": "Conversation", + "type": "P", # P = private channel (not a DM) + } + ) + except Exception as exc: + raise MattermostError(f"Channel creation failed: {exc}") from exc + + channel_id: str = channel["id"] + + for mm_user_id in mm_user_ids: + try: + self._driver.channels.add_user( + channel_id, options={"user_id": mm_user_id} + ) + except Exception as exc: + raise MattermostError( + f"Adding user {mm_user_id} to channel failed: {exc}" + ) from exc + + return channel_id + + def post_message(self, channel_id: str, message: str) -> str: + """ + Post a message to the given channel using the admin token. + Returns the opaque mattermost_post_id string. + """ + try: + post = self._driver.posts.create_post( + options={ + "channel_id": channel_id, + "message": message, + } + ) + return post["id"] + except Exception as exc: + raise MattermostError(f"Post creation failed: {exc}") from exc + + def get_posts( + self, + channel_id: str, + page: int = 0, + per_page: int = 20, + since: int | None = None, + ) -> list[dict]: + """ + Fetch posts from a channel and return them in chronological order. + + `since` is a Unix timestamp in **milliseconds**. When provided, + `page` is ignored and all posts after that timestamp are returned. + """ + params: dict = {"per_page": per_page} + if since is not None: + params["since"] = since + else: + params["page"] = page + + try: + data = self._driver.posts.get_posts_for_channel(channel_id, params=params) + except Exception as exc: + raise MattermostError(f"Fetching posts failed: {exc}") from exc + + order: list[str] = data.get("order", []) + posts_map: dict = data.get("posts", {}) + + # Mattermost returns newest-first; reverse to get chronological order. + return [posts_map[pid] for pid in reversed(order) if pid in posts_map] + + def get_latest_post_id(self, channel_id: str) -> str | None: + """ + Return the ID of the most recent post in the channel, or None if empty. + Used to derive `has_unread` against a stored last_read_mattermost_post_id. + """ + try: + data = self._driver.posts.get_posts_for_channel( + channel_id, params={"page": 0, "per_page": 1} + ) + except Exception as exc: + raise MattermostError(f"Fetching latest post failed: {exc}") from exc + + order: list[str] = data.get("order", []) + return order[0] if order else None diff --git a/apps/chat/integrations/mattermost/exceptions.py b/apps/chat/integrations/mattermost/exceptions.py new file mode 100644 index 0000000..2237877 --- /dev/null +++ b/apps/chat/integrations/mattermost/exceptions.py @@ -0,0 +1,2 @@ +class MattermostError(Exception): + """Raised when any Mattermost API call fails.""" diff --git a/apps/chat/migrations/0001_initial.py b/apps/chat/migrations/0001_initial.py new file mode 100644 index 0000000..e7cbb68 --- /dev/null +++ b/apps/chat/migrations/0001_initial.py @@ -0,0 +1,68 @@ +# Generated by Django 5.2.13 on 2026-06-10 12:09 + +import django.db.models.deletion +import uuid +from django.db import migrations, models + + +class Migration(migrations.Migration): + + initial = True + + dependencies = [ + ] + + operations = [ + migrations.CreateModel( + name='Conversation', + fields=[ + ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ('mattermost_channel_id', models.CharField(db_index=True, max_length=255, unique=True)), + ('type', models.CharField(choices=[('direct', 'Direct'), ('group', 'Group')], default='direct', max_length=20)), + ('created_at', models.DateTimeField(auto_now_add=True)), + ('updated_at', models.DateTimeField(auto_now=True)), + ], + options={ + 'indexes': [models.Index(fields=['type'], name='chat_conver_type_857bd7_idx'), models.Index(fields=['created_at'], name='chat_conver_created_656b50_idx')], + }, + ), + migrations.CreateModel( + name='MattermostAccountMapping', + fields=[ + ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ('user_id', models.UUIDField(db_index=True, unique=True)), + ('mattermost_user_id', models.CharField(db_index=True, max_length=255, unique=True)), + ('created_at', models.DateTimeField(auto_now_add=True)), + ], + options={ + 'indexes': [models.Index(fields=['mattermost_user_id'], name='chat_matter_matterm_e23c84_idx')], + }, + ), + migrations.CreateModel( + name='ConversationParticipant', + fields=[ + ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ('user_id', models.UUIDField()), + ('joined_at', models.DateTimeField(auto_now_add=True)), + ('conversation', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='participants', to='chat.conversation')), + ], + options={ + 'indexes': [models.Index(fields=['user_id'], name='chat_conver_user_id_697713_idx'), models.Index(fields=['conversation', 'user_id'], name='chat_conver_convers_59b4dd_idx')], + 'constraints': [models.UniqueConstraint(fields=('conversation', 'user_id'), name='unique_participant_per_conversation')], + }, + ), + migrations.CreateModel( + name='ConversationReadState', + fields=[ + ('id', models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ('user_id', models.UUIDField()), + ('last_read_mattermost_post_id', models.CharField(blank=True, max_length=255, null=True)), + ('updated_at', models.DateTimeField(auto_now=True)), + ('conversation', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='read_states', to='chat.conversation')), + ], + options={ + 'indexes': [models.Index(fields=['user_id'], name='chat_conver_user_id_e7609f_idx'), models.Index(fields=['conversation', 'user_id'], name='chat_conver_convers_9555aa_idx')], + 'constraints': [models.UniqueConstraint(fields=('conversation', 'user_id'), name='unique_read_state_per_user')], + }, + ), + ] diff --git a/apps/chat/migrations/0002_remove_conversationparticipant_unique_participant_per_conversation_and_more.py b/apps/chat/migrations/0002_remove_conversationparticipant_unique_participant_per_conversation_and_more.py new file mode 100644 index 0000000..ae7a301 --- /dev/null +++ b/apps/chat/migrations/0002_remove_conversationparticipant_unique_participant_per_conversation_and_more.py @@ -0,0 +1,96 @@ +# Generated by Django 5.2.13 on 2026-07-19 08:24 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('chat', '0001_initial'), + ] + + operations = [ + migrations.RemoveConstraint( + model_name='conversationparticipant', + name='unique_participant_per_conversation', + ), + migrations.RemoveConstraint( + model_name='conversationreadstate', + name='unique_read_state_per_user', + ), + migrations.RemoveIndex( + model_name='conversationparticipant', + name='chat_conver_user_id_697713_idx', + ), + migrations.RemoveIndex( + model_name='conversationparticipant', + name='chat_conver_convers_59b4dd_idx', + ), + migrations.RemoveIndex( + model_name='conversationreadstate', + name='chat_conver_user_id_e7609f_idx', + ), + migrations.RemoveIndex( + model_name='conversationreadstate', + name='chat_conver_convers_9555aa_idx', + ), + migrations.RenameField( + model_name='conversation', + old_name='id', + new_name='uuid', + ), + migrations.RenameField( + model_name='conversationparticipant', + old_name='user_id', + new_name='user_uuid', + ), + migrations.RenameField( + model_name='conversationparticipant', + old_name='id', + new_name='uuid', + ), + migrations.RenameField( + model_name='conversationreadstate', + old_name='user_id', + new_name='user_uuid', + ), + migrations.RenameField( + model_name='conversationreadstate', + old_name='id', + new_name='uuid', + ), + migrations.RenameField( + model_name='mattermostaccountmapping', + old_name='user_id', + new_name='user_uuid', + ), + migrations.RenameField( + model_name='mattermostaccountmapping', + old_name='id', + new_name='uuid', + ), + migrations.AddIndex( + model_name='conversationparticipant', + index=models.Index(fields=['user_uuid'], name='chat_conver_user_uu_020064_idx'), + ), + migrations.AddIndex( + model_name='conversationparticipant', + index=models.Index(fields=['conversation', 'user_uuid'], name='chat_conver_convers_5e51b5_idx'), + ), + migrations.AddIndex( + model_name='conversationreadstate', + index=models.Index(fields=['user_uuid'], name='chat_conver_user_uu_8e7f04_idx'), + ), + migrations.AddIndex( + model_name='conversationreadstate', + index=models.Index(fields=['conversation', 'user_uuid'], name='chat_conver_convers_d7112b_idx'), + ), + migrations.AddConstraint( + model_name='conversationparticipant', + constraint=models.UniqueConstraint(fields=('conversation', 'user_uuid'), name='unique_participant_per_conversation'), + ), + migrations.AddConstraint( + model_name='conversationreadstate', + constraint=models.UniqueConstraint(fields=('conversation', 'user_uuid'), name='unique_read_state_per_user'), + ), + ] diff --git a/apps/chat/migrations/0003_conversation_closed_by_uuid_conversation_status_and_more.py b/apps/chat/migrations/0003_conversation_closed_by_uuid_conversation_status_and_more.py new file mode 100644 index 0000000..41abdd4 --- /dev/null +++ b/apps/chat/migrations/0003_conversation_closed_by_uuid_conversation_status_and_more.py @@ -0,0 +1,27 @@ +# Generated by Django 5.2.13 on 2026-07-19 09:05 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('chat', '0002_remove_conversationparticipant_unique_participant_per_conversation_and_more'), + ] + + operations = [ + migrations.AddField( + model_name='conversation', + name='closed_by_uuid', + field=models.UUIDField(blank=True, null=True), + ), + migrations.AddField( + model_name='conversation', + name='status', + field=models.CharField(choices=[('open', 'Open'), ('closed', 'Closed')], default='open', max_length=20), + ), + migrations.AddIndex( + model_name='conversation', + index=models.Index(fields=['status'], name='chat_conver_status_9f0685_idx'), + ), + ] diff --git a/apps/chat/models.py b/apps/chat/models.py deleted file mode 100644 index 71a8362..0000000 --- a/apps/chat/models.py +++ /dev/null @@ -1,3 +0,0 @@ -from django.db import models - -# Create your models here. diff --git a/apps/chat/models/__init__.py b/apps/chat/models/__init__.py new file mode 100644 index 0000000..873fa0d --- /dev/null +++ b/apps/chat/models/__init__.py @@ -0,0 +1,13 @@ +from apps.chat.models.conversation import Conversation, ConversationStatus, ConversationType +from apps.chat.models.mapping import MattermostAccountMapping +from apps.chat.models.participant import ConversationParticipant +from apps.chat.models.read_state import ConversationReadState + +__all__ = [ + "Conversation", + "ConversationType", + "ConversationStatus", + "ConversationParticipant", + "MattermostAccountMapping", + "ConversationReadState", +] diff --git a/apps/chat/models/conversation.py b/apps/chat/models/conversation.py new file mode 100644 index 0000000..2f67bc1 --- /dev/null +++ b/apps/chat/models/conversation.py @@ -0,0 +1,43 @@ +import uuid + +from django.db import models + + +class ConversationType(models.TextChoices): + DIRECT = "direct", "Direct" + GROUP = "group", "Group" + + +class ConversationStatus(models.TextChoices): + OPEN = "open", "Open" + CLOSED = "closed", "Closed" + + +class Conversation(models.Model): + uuid = models.UUIDField(primary_key=True, default=uuid.uuid4, editable=False) + # Internal only — never exposed through any public API or serializer + mattermost_channel_id = models.CharField(max_length=255, unique=True, db_index=True) + type = models.CharField( + max_length=20, + choices=ConversationType.choices, + default=ConversationType.DIRECT, + ) + status = models.CharField( + max_length=20, + choices=ConversationStatus.choices, + default=ConversationStatus.OPEN, + ) + # Set only while status is CLOSED; identifies who is allowed to reopen it. + closed_by_uuid = models.UUIDField(null=True, blank=True) + created_at = models.DateTimeField(auto_now_add=True) + updated_at = models.DateTimeField(auto_now=True) + + class Meta: + indexes = [ + models.Index(fields=["type"]), + models.Index(fields=["created_at"]), + models.Index(fields=["status"]), + ] + + def __str__(self): + return str(self.uuid) diff --git a/apps/chat/models/mapping.py b/apps/chat/models/mapping.py new file mode 100644 index 0000000..4254b97 --- /dev/null +++ b/apps/chat/models/mapping.py @@ -0,0 +1,18 @@ +import uuid + +from django.db import models + + +class MattermostAccountMapping(models.Model): + uuid = models.UUIDField(primary_key=True, default=uuid.uuid4, editable=False) + user_uuid = models.UUIDField(unique=True, db_index=True) + mattermost_user_id = models.CharField(max_length=255, unique=True, db_index=True) + created_at = models.DateTimeField(auto_now_add=True) + + class Meta: + indexes = [ + models.Index(fields=["mattermost_user_id"]), + ] + + def __str__(self): + return f"{self.user_uuid} -> {self.mattermost_user_id}" diff --git a/apps/chat/models/participant.py b/apps/chat/models/participant.py new file mode 100644 index 0000000..be150dc --- /dev/null +++ b/apps/chat/models/participant.py @@ -0,0 +1,31 @@ +import uuid + +from django.db import models + +from apps.chat.models.conversation import Conversation + + +class ConversationParticipant(models.Model): + uuid = models.UUIDField(primary_key=True, default=uuid.uuid4, editable=False) + conversation = models.ForeignKey( + Conversation, + related_name="participants", + on_delete=models.CASCADE, + ) + user_uuid = models.UUIDField() + joined_at = models.DateTimeField(auto_now_add=True) + + class Meta: + constraints = [ + models.UniqueConstraint( + fields=["conversation", "user_uuid"], + name="unique_participant_per_conversation", + ) + ] + indexes = [ + models.Index(fields=["user_uuid"]), + models.Index(fields=["conversation", "user_uuid"]), + ] + + def __str__(self): + return f"{self.user_uuid} in {self.conversation_id}" diff --git a/apps/chat/models/read_state.py b/apps/chat/models/read_state.py new file mode 100644 index 0000000..e28a356 --- /dev/null +++ b/apps/chat/models/read_state.py @@ -0,0 +1,32 @@ +import uuid + +from django.db import models + +from apps.chat.models.conversation import Conversation + + +class ConversationReadState(models.Model): + uuid = models.UUIDField(primary_key=True, default=uuid.uuid4, editable=False) + conversation = models.ForeignKey( + Conversation, + related_name="read_states", + on_delete=models.CASCADE, + ) + user_uuid = models.UUIDField() + last_read_mattermost_post_id = models.CharField(max_length=255, null=True, blank=True) + updated_at = models.DateTimeField(auto_now=True) + + class Meta: + constraints = [ + models.UniqueConstraint( + fields=["conversation", "user_uuid"], + name="unique_read_state_per_user", + ) + ] + indexes = [ + models.Index(fields=["user_uuid"]), + models.Index(fields=["conversation", "user_uuid"]), + ] + + def __str__(self): + return f"{self.user_uuid} read state for {self.conversation_id}" diff --git a/apps/chat/serializers/__init__.py b/apps/chat/serializers/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/serializers/conversations.py b/apps/chat/serializers/conversations.py new file mode 100644 index 0000000..d757ecf --- /dev/null +++ b/apps/chat/serializers/conversations.py @@ -0,0 +1,40 @@ +from drf_spectacular.utils import extend_schema_field +from rest_framework import serializers + +from apps.chat.models import Conversation, ConversationParticipant + + +class CreateConversationSerializer(serializers.Serializer): + user_1_uuid = serializers.UUIDField() + user_2_uuid = serializers.UUIDField() + + +class ConversationUserActionSerializer(serializers.Serializer): + user_uuid = serializers.UUIDField() + + +class ConversationSerializer(serializers.ModelSerializer): + participants = serializers.SerializerMethodField() + + class Meta: + model = Conversation + fields = ["uuid", "type", "status", "closed_by_uuid", "created_at", "participants"] + + @extend_schema_field(serializers.ListField(child=serializers.UUIDField())) + def get_participants(self, obj): + return list( + ConversationParticipant.objects.filter(conversation_id=obj.uuid) + .order_by("joined_at") + .values_list("user_uuid", flat=True) + ) + + +class ConversationListSerializer(serializers.Serializer): + uuid = serializers.UUIDField() + type = serializers.CharField() + status = serializers.CharField() + closed_by_uuid = serializers.UUIDField(allow_null=True) + created_at = serializers.DateTimeField() + participants = serializers.ListField(child=serializers.UUIDField()) + has_unread = serializers.BooleanField() + last_message = serializers.DictField(allow_null=True, required=False) diff --git a/apps/chat/serializers/messages.py b/apps/chat/serializers/messages.py new file mode 100644 index 0000000..67bc47b --- /dev/null +++ b/apps/chat/serializers/messages.py @@ -0,0 +1,17 @@ +from rest_framework import serializers + + +class SendMessageSerializer(serializers.Serializer): + sender_uuid = serializers.UUIDField() + message_type = serializers.ChoiceField(choices=["text", "image", "file"], default="text") + text = serializers.CharField(required=False, allow_null=True, allow_blank=True) + file = serializers.FileField(required=False, allow_null=True) + + +class MessageSerializer(serializers.Serializer): + post_id = serializers.CharField() + sender_uuid = serializers.UUIDField(allow_null=True) + message_type = serializers.CharField(default="text") + text = serializers.CharField(allow_null=True, required=False) + url = serializers.CharField(allow_null=True, required=False) + created_at = serializers.DateTimeField(allow_null=True, required=False) diff --git a/apps/chat/serializers/read_state.py b/apps/chat/serializers/read_state.py new file mode 100644 index 0000000..c8c9137 --- /dev/null +++ b/apps/chat/serializers/read_state.py @@ -0,0 +1,12 @@ +from rest_framework import serializers + + +class MarkReadSerializer(serializers.Serializer): + user_uuid = serializers.UUIDField() + post_id = serializers.CharField() + + +class ReadStateSerializer(serializers.Serializer): + conversation_uuid = serializers.UUIDField() + user_uuid = serializers.UUIDField() + has_unread = serializers.BooleanField() diff --git a/apps/chat/services/__init__.py b/apps/chat/services/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/services/account.py b/apps/chat/services/account.py new file mode 100644 index 0000000..163951f --- /dev/null +++ b/apps/chat/services/account.py @@ -0,0 +1,7 @@ +from uuid import UUID + + +class AccountService: + def validate_user(self, user_uuid: UUID) -> bool: + # Future: call external Account Service HTTP API + return True diff --git a/apps/chat/services/conversation.py b/apps/chat/services/conversation.py new file mode 100644 index 0000000..b11fe9b --- /dev/null +++ b/apps/chat/services/conversation.py @@ -0,0 +1,91 @@ +from uuid import UUID + +from django.db import transaction + +from apps.chat.integrations.mattermost.client import MattermostClient +from apps.chat.models import ( + Conversation, + ConversationParticipant, + ConversationStatus, + MattermostAccountMapping, +) +from apps.chat.services.account import AccountService + + +class ConversationService: + def __init__( + self, + account_service: AccountService | None = None, + mattermost_client: MattermostClient | None = None, + ): + self._account = account_service or AccountService() + self._mm = mattermost_client or MattermostClient() + + def create(self, user_1_uuid: UUID, user_2_uuid: UUID) -> Conversation: + if not self._account.validate_user(user_1_uuid): + raise ValueError(f"User {user_1_uuid} is not valid") + if not self._account.validate_user(user_2_uuid): + raise ValueError(f"User {user_2_uuid} is not valid") + + mm_user_1 = self._mm.get_or_create_user(user_1_uuid) + mm_user_2 = self._mm.get_or_create_user(user_2_uuid) + + MattermostAccountMapping.objects.get_or_create( + user_uuid=user_1_uuid, + defaults={"mattermost_user_id": mm_user_1}, + ) + MattermostAccountMapping.objects.get_or_create( + user_uuid=user_2_uuid, + defaults={"mattermost_user_id": mm_user_2}, + ) + + channel_id = self._mm.create_private_channel([mm_user_1, mm_user_2]) + + with transaction.atomic(): + conversation = Conversation.objects.create(mattermost_channel_id=channel_id) + ConversationParticipant.objects.create( + conversation=conversation, user_uuid=user_1_uuid + ) + ConversationParticipant.objects.create( + conversation=conversation, user_uuid=user_2_uuid + ) + + return conversation + + def list_for_user( + self, user_uuid: UUID, page: int, page_size: int + ) -> list[Conversation]: + offset = page * page_size + conversation_uuids = ConversationParticipant.objects.filter( + user_uuid=user_uuid + ).values_list("conversation_id", flat=True) + return list( + Conversation.objects.filter(uuid__in=conversation_uuids).order_by( + "-created_at" + )[offset : offset + page_size] + ) + + def close(self, conversation_uuid: UUID, user_uuid: UUID) -> Conversation: + conversation = Conversation.objects.get(uuid=conversation_uuid) + if conversation.status == ConversationStatus.CLOSED: + return conversation + + conversation.status = ConversationStatus.CLOSED + conversation.closed_by_uuid = user_uuid + conversation.save(update_fields=["status", "closed_by_uuid", "updated_at"]) + return conversation + + def reopen(self, conversation_uuid: UUID, user_uuid: UUID) -> Conversation: + conversation = Conversation.objects.get(uuid=conversation_uuid) + if conversation.status == ConversationStatus.OPEN: + return conversation + + if conversation.closed_by_uuid != user_uuid: + raise PermissionError( + "Only the user who closed this conversation can reopen it" + ) + + conversation.status = ConversationStatus.OPEN + conversation.closed_by_uuid = None + conversation.save(update_fields=["status", "closed_by_uuid", "updated_at"]) + return conversation diff --git a/apps/chat/services/message.py b/apps/chat/services/message.py new file mode 100644 index 0000000..e2c7236 --- /dev/null +++ b/apps/chat/services/message.py @@ -0,0 +1,128 @@ +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 + + +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, + file=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") + + file_url: str | None = None + if file is not None: + filename = getattr(file, "name", f"{sender_uuid}") + file_url = self._storage.upload_file(file, filename) + + if message_type == "text": + mm_message = text or "" + else: + mm_message = file_url or 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, "file_url": file_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": file_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] + + @staticmethod + def _normalize_post(post: dict, mappings: dict) -> dict: + mm_uid = post.get("user_id") + sender_uuid = mappings.get(mm_uid) + msg = post.get("message", "") + is_url = msg.startswith("http") + created_ms = post.get("create_at") + created_at = ( + datetime.fromtimestamp(created_ms / 1000, tz=timezone.utc) + if created_ms + else None + ) + return { + "post_id": post.get("id"), + "sender_uuid": sender_uuid, + "message_type": "file" if is_url else "text", + "text": None if is_url else msg, + "url": msg if is_url else None, + "created_at": created_at, + } diff --git a/apps/chat/services/read_state.py b/apps/chat/services/read_state.py new file mode 100644 index 0000000..8cb2a86 --- /dev/null +++ b/apps/chat/services/read_state.py @@ -0,0 +1,48 @@ +from uuid import UUID + +from django.core.cache import cache + +from apps.chat.integrations.mattermost.client import MattermostClient +from apps.chat.models import Conversation, ConversationReadState + +_LATEST_POST_CACHE_TTL = 60 + + +class ReadStateService: + def __init__(self, mattermost_client: MattermostClient | None = None): + self._mm = mattermost_client or MattermostClient() + + def mark_read( + self, conversation_uuid, user_uuid: UUID, post_id: str + ) -> ConversationReadState: + read_state, _ = ConversationReadState.objects.update_or_create( + conversation_id=conversation_uuid, + user_uuid=user_uuid, + defaults={"last_read_mattermost_post_id": post_id}, + ) + return read_state + + def has_unread(self, conversation_uuid, user_uuid: UUID) -> bool: + try: + read_state = ConversationReadState.objects.get( + conversation_id=conversation_uuid, user_uuid=user_uuid + ) + except ConversationReadState.DoesNotExist: + return True + + last_read = read_state.last_read_mattermost_post_id + + cache_key = f"latest_post:{conversation_uuid}" + latest_post_id = cache.get(cache_key) + + if latest_post_id is None: + conversation = Conversation.objects.get(uuid=conversation_uuid) + latest_post_id = self._mm.get_latest_post_id( + conversation.mattermost_channel_id + ) + cache.set(cache_key, latest_post_id or "", _LATEST_POST_CACHE_TTL) + + if not latest_post_id: + return False + + return latest_post_id != last_read diff --git a/apps/chat/services/realtime.py b/apps/chat/services/realtime.py new file mode 100644 index 0000000..1d30ea1 --- /dev/null +++ b/apps/chat/services/realtime.py @@ -0,0 +1,12 @@ +from asgiref.sync import async_to_sync +from channels.layers import get_channel_layer + + +class RealtimeService: + def publish_message(self, chat_uuid, payload: dict) -> None: + channel_layer = get_channel_layer() + group_name = f"chat_{chat_uuid}" + async_to_sync(channel_layer.group_send)( + group_name, + {"type": "chat.message", **payload}, + ) diff --git a/apps/chat/services/storage.py b/apps/chat/services/storage.py new file mode 100644 index 0000000..eff1738 --- /dev/null +++ b/apps/chat/services/storage.py @@ -0,0 +1,42 @@ +import io +import mimetypes +from uuid import uuid4 + +from django.conf import settings +from minio import Minio + + +class StorageService: + def __init__(self): + endpoint = getattr(settings, "MINIO_ENDPOINT", None) + if not endpoint: + raise ValueError("MINIO_ENDPOINT is not configured") + + secure = getattr(settings, "MINIO_SECURE", False) + self._client = Minio( + endpoint, + access_key=getattr(settings, "MINIO_ACCESS_KEY", None), + secret_key=getattr(settings, "MINIO_SECRET_KEY", None), + secure=secure, + ) + self._bucket = getattr(settings, "MINIO_BUCKET_CHAT", "chat") + scheme = "https" if secure else "http" + self._base_url = f"{scheme}://{endpoint}" + + def upload_file(self, file, filename: str) -> str: + data = file.read() if hasattr(file, "read") else file + size = len(data) + + content_type, _ = mimetypes.guess_type(filename) + if not content_type: + content_type = "application/octet-stream" + + object_name = f"{uuid4().hex}/{filename}" + self._client.put_object( + self._bucket, + object_name, + io.BytesIO(data), + size, + content_type=content_type, + ) + return f"{self._base_url}/{self._bucket}/{object_name}" diff --git a/apps/chat/tests.py b/apps/chat/tests.py deleted file mode 100644 index 7ce503c..0000000 --- a/apps/chat/tests.py +++ /dev/null @@ -1,3 +0,0 @@ -from django.test import TestCase - -# Create your tests here. diff --git a/apps/chat/tests/__init__.py b/apps/chat/tests/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/tests/test_conversation_close.py b/apps/chat/tests/test_conversation_close.py new file mode 100644 index 0000000..4092c29 --- /dev/null +++ b/apps/chat/tests/test_conversation_close.py @@ -0,0 +1,112 @@ +import uuid +from unittest.mock import Mock + +import pytest + +from apps.chat.exceptions import ConversationClosedError +from apps.chat.models import Conversation, ConversationStatus +from apps.chat.services.conversation import ConversationService +from apps.chat.services.message import MessageService + + +def _make_conv(status=ConversationStatus.OPEN, closed_by_uuid=None, channel_id="ch-close"): + return Conversation.objects.create( + mattermost_channel_id=channel_id, + status=status, + closed_by_uuid=closed_by_uuid, + ) + + +def _svc(): + return ConversationService(mattermost_client=Mock()) + + +@pytest.mark.django_db +def test_close_sets_status_and_closed_by(): + conv = _make_conv() + closer = uuid.uuid4() + + result = _svc().close(conv.uuid, closer) + + assert result.status == ConversationStatus.CLOSED + assert result.closed_by_uuid == closer + + conv.refresh_from_db() + assert conv.status == ConversationStatus.CLOSED + assert conv.closed_by_uuid == closer + + +@pytest.mark.django_db +def test_close_is_idempotent(): + closer = uuid.uuid4() + other = uuid.uuid4() + conv = _make_conv(status=ConversationStatus.CLOSED, closed_by_uuid=closer) + + result = _svc().close(conv.uuid, other) + + assert result.status == ConversationStatus.CLOSED + assert result.closed_by_uuid == closer # unchanged — original closer preserved + + +@pytest.mark.django_db +def test_reopen_by_closer_succeeds(): + closer = uuid.uuid4() + conv = _make_conv(status=ConversationStatus.CLOSED, closed_by_uuid=closer) + + result = _svc().reopen(conv.uuid, closer) + + assert result.status == ConversationStatus.OPEN + assert result.closed_by_uuid is None + + +@pytest.mark.django_db +def test_reopen_by_other_user_raises_permission_error(): + closer = uuid.uuid4() + other = uuid.uuid4() + conv = _make_conv(status=ConversationStatus.CLOSED, closed_by_uuid=closer) + + with pytest.raises(PermissionError): + _svc().reopen(conv.uuid, other) + + conv.refresh_from_db() + assert conv.status == ConversationStatus.CLOSED + + +@pytest.mark.django_db +def test_reopen_is_idempotent_when_already_open(): + conv = _make_conv(status=ConversationStatus.OPEN) + someone = uuid.uuid4() + + result = _svc().reopen(conv.uuid, someone) + + assert result.status == ConversationStatus.OPEN + + +@pytest.mark.django_db +def test_send_message_to_closed_conversation_raises(): + closer = uuid.uuid4() + conv = _make_conv(status=ConversationStatus.CLOSED, closed_by_uuid=closer) + sender_uuid = uuid.uuid4() + + svc = MessageService(mattermost_client=Mock(), storage_service=Mock(), publishers=[]) + + with pytest.raises(ConversationClosedError): + svc.send(conv.uuid, sender_uuid, "text", text="hello") + + +@pytest.mark.django_db +def test_send_message_after_reopen_succeeds(): + closer = uuid.uuid4() + conv = _make_conv(status=ConversationStatus.CLOSED, closed_by_uuid=closer) + sender_uuid = uuid.uuid4() + + _svc().reopen(conv.uuid, closer) + + mm_client = Mock() + mm_client.post_message.return_value = "post-1" + svc = MessageService(mattermost_client=mm_client, storage_service=Mock(), publishers=[]) + + result = svc.send(conv.uuid, sender_uuid, "text", text="hello again") + + assert result["post_id"] == "post-1" + mm_client.post_message.assert_called_once_with("ch-close", "hello again") diff --git a/apps/chat/tests/test_conversations.py b/apps/chat/tests/test_conversations.py new file mode 100644 index 0000000..278a14e --- /dev/null +++ b/apps/chat/tests/test_conversations.py @@ -0,0 +1,73 @@ +import uuid +from unittest.mock import Mock + +import pytest + +from apps.chat.models import ( + Conversation, + ConversationParticipant, + MattermostAccountMapping, +) +from apps.chat.services.conversation import ConversationService + + +def _make_mm_client(mm_user_ids=("mm-u1", "mm-u2"), channel_id="mm-ch-1"): + client = Mock() + client.get_or_create_user.side_effect = list(mm_user_ids) + client.create_private_channel.return_value = channel_id + return client + + +@pytest.mark.django_db +def test_create_conversation_success(): + user_1 = uuid.uuid4() + user_2 = uuid.uuid4() + + svc = ConversationService(mattermost_client=_make_mm_client()) + conv = svc.create(user_1, user_2) + + assert Conversation.objects.filter(uuid=conv.uuid).exists() + + participants = ConversationParticipant.objects.filter(conversation=conv) + assert participants.count() == 2 + participant_uuids = set(participants.values_list("user_uuid", flat=True)) + assert user_1 in participant_uuids + assert user_2 in participant_uuids + + +@pytest.mark.django_db +def test_create_conversation_creates_mattermost_mapping(): + user_1 = uuid.uuid4() + user_2 = uuid.uuid4() + + mm_client = _make_mm_client(mm_user_ids=("mm-user-aaa", "mm-user-bbb")) + svc = ConversationService(mattermost_client=mm_client) + svc.create(user_1, user_2) + + mapping_1 = MattermostAccountMapping.objects.get(user_uuid=user_1) + mapping_2 = MattermostAccountMapping.objects.get(user_uuid=user_2) + + assert mapping_1.mattermost_user_id == "mm-user-aaa" + assert mapping_2.mattermost_user_id == "mm-user-bbb" + + +@pytest.mark.django_db +def test_list_for_user_returns_only_user_conversations(): + target_user = uuid.uuid4() + other_user = uuid.uuid4() + + conv_with_user = Conversation.objects.create(mattermost_channel_id="ch-target") + ConversationParticipant.objects.create( + conversation=conv_with_user, user_uuid=target_user + ) + + conv_without_user = Conversation.objects.create(mattermost_channel_id="ch-other") + ConversationParticipant.objects.create( + conversation=conv_without_user, user_uuid=other_user + ) + + svc = ConversationService(mattermost_client=Mock()) + result = svc.list_for_user(target_user, page=0, page_size=20) + + assert len(result) == 1 + assert result[0].uuid == conv_with_user.uuid diff --git a/apps/chat/tests/test_messages.py b/apps/chat/tests/test_messages.py new file mode 100644 index 0000000..294d606 --- /dev/null +++ b/apps/chat/tests/test_messages.py @@ -0,0 +1,116 @@ +import uuid +from unittest.mock import Mock + +import pytest + +from apps.chat.models import Conversation, MattermostAccountMapping +from apps.chat.services.message import MessageService + + +def _make_service(*, mm_post_id="post-1", storage_url=None, publishers=None): + """Return a MessageService with all external deps mocked.""" + mm_client = Mock() + mm_client.post_message.return_value = mm_post_id + + storage = Mock() + if storage_url: + storage.upload_file.return_value = storage_url + + return ( + MessageService( + mattermost_client=mm_client, + storage_service=storage, + publishers=publishers if publishers is not None else [], + ), + mm_client, + storage, + ) + + +@pytest.mark.django_db +def test_send_text_message(): + conv = Conversation.objects.create(mattermost_channel_id="ch-send-text") + sender_uuid = uuid.uuid4() + + ws_publisher = Mock() + svc, mm_client, _ = _make_service(mm_post_id="post-abc", publishers=[ws_publisher]) + + result = svc.send(conv.uuid, sender_uuid, "text", text="Hello world") + + mm_client.post_message.assert_called_once_with("ch-send-text", "Hello world") + ws_publisher.publish.assert_called_once() + + event = ws_publisher.publish.call_args[0][0] + assert event.post_id == "post-abc" + assert event.sender_uuid == sender_uuid + assert event.message_type == "text" + + assert result["post_id"] == "post-abc" + assert result["text"] == "Hello world" + + +@pytest.mark.django_db +def test_send_image_message(): + conv = Conversation.objects.create(mattermost_channel_id="ch-send-image") + sender_uuid = uuid.uuid4() + + file_url = "http://minio.local/chat/abc123/photo.jpg" + svc, mm_client, storage = _make_service(storage_url=file_url, publishers=[]) + + fake_file = Mock() + fake_file.name = "photo.jpg" + + result = svc.send(conv.uuid, sender_uuid, "image", file=fake_file) + + storage.upload_file.assert_called_once_with(fake_file, "photo.jpg") + mm_client.post_message.assert_called_once_with("ch-send-image", file_url) + assert result["url"] == file_url + + +@pytest.mark.django_db +def test_list_messages_normalized(): + conv = Conversation.objects.create(mattermost_channel_id="ch-list") + sender_uuid = uuid.uuid4() + mm_user_id = "mm-user-xyz" + + MattermostAccountMapping.objects.create( + user_uuid=sender_uuid, mattermost_user_id=mm_user_id + ) + + raw_posts = [ + { + "id": "p1", + "user_id": mm_user_id, + "message": "First message", + "create_at": 1700000000000, + }, + { + "id": "p2", + "user_id": mm_user_id, + "message": "http://minio.local/chat/file.pdf", + "create_at": 1700000001000, + }, + ] + + mm_client = Mock() + mm_client.get_posts.return_value = raw_posts + svc, _, _ = _make_service(publishers=[]) + svc._mm = mm_client # inject after construction to keep _make_service simple + + messages = svc.list_messages(conv.uuid, page=0, per_page=20) + + assert len(messages) == 2 + + text_msg = messages[0] + assert text_msg["post_id"] == "p1" + assert text_msg["sender_uuid"] == sender_uuid + assert text_msg["message_type"] == "text" + assert text_msg["text"] == "First message" + assert text_msg["url"] is None + + file_msg = messages[1] + assert file_msg["post_id"] == "p2" + assert file_msg["sender_uuid"] == sender_uuid + assert file_msg["message_type"] == "file" + assert file_msg["url"] == "http://minio.local/chat/file.pdf" + assert file_msg["text"] is None diff --git a/apps/chat/tests/test_read_state.py b/apps/chat/tests/test_read_state.py new file mode 100644 index 0000000..93c8b18 --- /dev/null +++ b/apps/chat/tests/test_read_state.py @@ -0,0 +1,86 @@ +import uuid +from unittest.mock import Mock, patch + +import pytest + +from apps.chat.models import Conversation, ConversationReadState +from apps.chat.services.read_state import ReadStateService + + +def _make_conv(channel_id="ch-read"): + return Conversation.objects.create(mattermost_channel_id=channel_id) + + +@pytest.mark.django_db +def test_mark_read_creates_read_state(): + conv = _make_conv() + user_uuid = uuid.uuid4() + + svc = ReadStateService(mattermost_client=Mock()) + read_state = svc.mark_read(conv.uuid, user_uuid, "post-111") + + assert read_state.conversation_id == conv.uuid + assert read_state.user_uuid == user_uuid + assert read_state.last_read_mattermost_post_id == "post-111" + assert ConversationReadState.objects.filter( + conversation=conv, user_uuid=user_uuid + ).count() == 1 + + +@pytest.mark.django_db +def test_mark_read_updates_existing(): + conv = _make_conv() + user_uuid = uuid.uuid4() + + svc = ReadStateService(mattermost_client=Mock()) + svc.mark_read(conv.uuid, user_uuid, "post-old") + svc.mark_read(conv.uuid, user_uuid, "post-new") + + rows = ConversationReadState.objects.filter(conversation=conv, user_uuid=user_uuid) + assert rows.count() == 1 + assert rows.first().last_read_mattermost_post_id == "post-new" + + +@pytest.mark.django_db +@patch("apps.chat.services.read_state.cache") +def test_has_unread_true(mock_cache): + mock_cache.get.return_value = None # force cache miss → MM lookup + + conv = _make_conv() + user_uuid = uuid.uuid4() + ConversationReadState.objects.create( + conversation=conv, + user_uuid=user_uuid, + last_read_mattermost_post_id="post-old", + ) + + mm_client = Mock() + mm_client.get_latest_post_id.return_value = "post-new" + + svc = ReadStateService(mattermost_client=mm_client) + result = svc.has_unread(conv.uuid, user_uuid) + + assert result is True + mm_client.get_latest_post_id.assert_called_once_with(conv.mattermost_channel_id) + + +@pytest.mark.django_db +@patch("apps.chat.services.read_state.cache") +def test_has_unread_false(mock_cache): + mock_cache.get.return_value = None # force cache miss → MM lookup + + conv = _make_conv() + user_uuid = uuid.uuid4() + ConversationReadState.objects.create( + conversation=conv, + user_uuid=user_uuid, + last_read_mattermost_post_id="post-current", + ) + + mm_client = Mock() + mm_client.get_latest_post_id.return_value = "post-current" + + svc = ReadStateService(mattermost_client=mm_client) + result = svc.has_unread(conv.uuid, user_uuid) + + assert result is False diff --git a/apps/chat/urls.py b/apps/chat/urls.py new file mode 100644 index 0000000..1571a28 --- /dev/null +++ b/apps/chat/urls.py @@ -0,0 +1,20 @@ +from django.urls import path + +from apps.chat.views.conversations import ( + ConversationCloseView, + ConversationCreateView, + ConversationReopenView, + UserConversationListView, +) +from apps.chat.views.messages import MessageView +from apps.chat.views.read_state import ChatEventsView, ReadStateView + +urlpatterns = [ + path("api/chats/", ConversationCreateView.as_view(), name="chat-create"), + path("api/users//chats/", UserConversationListView.as_view(), name="user-chat-list"), + path("api/chats//messages/", MessageView.as_view(), name="chat-messages"), + path("api/chats//read/", ReadStateView.as_view(), name="chat-read"), + path("api/chats//events/", ChatEventsView.as_view(), name="chat-events"), + path("api/chats//close/", ConversationCloseView.as_view(), name="chat-close"), + path("api/chats//reopen/", ConversationReopenView.as_view(), name="chat-reopen"), +] diff --git a/apps/chat/views.py b/apps/chat/views.py deleted file mode 100644 index 91ea44a..0000000 --- a/apps/chat/views.py +++ /dev/null @@ -1,3 +0,0 @@ -from django.shortcuts import render - -# Create your views here. diff --git a/apps/chat/views/__init__.py b/apps/chat/views/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/views/conversations.py b/apps/chat/views/conversations.py new file mode 100644 index 0000000..f1dc8e1 --- /dev/null +++ b/apps/chat/views/conversations.py @@ -0,0 +1,138 @@ +from uuid import UUID + +from drf_spectacular.utils import OpenApiParameter, extend_schema +from rest_framework import status +from rest_framework.response import Response +from rest_framework.views import APIView + +from apps.chat.models import ConversationParticipant +from apps.chat.serializers.conversations import ( + ConversationListSerializer, + ConversationSerializer, + ConversationUserActionSerializer, + CreateConversationSerializer, +) +from apps.chat.services.conversation import ConversationService +from apps.chat.services.message import MessageService +from apps.chat.services.read_state import ReadStateService + + +class ConversationCreateView(APIView): + authentication_classes = [] + permission_classes = [] + + @extend_schema( + request=CreateConversationSerializer, + responses={201: ConversationSerializer}, + ) + def post(self, request): + serializer = CreateConversationSerializer(data=request.data) + serializer.is_valid(raise_exception=True) + + conversation = ConversationService().create( + serializer.validated_data["user_1_uuid"], + serializer.validated_data["user_2_uuid"], + ) + + return Response( + ConversationSerializer(conversation).data, + status=status.HTTP_201_CREATED, + ) + + +class UserConversationListView(APIView): + authentication_classes = [] + permission_classes = [] + + @extend_schema( + parameters=[ + OpenApiParameter("page", int, description="0-based page number", default=0), + OpenApiParameter("page_size", int, description="Results per page", default=20), + ], + responses={200: ConversationListSerializer(many=True)}, + ) + def get(self, request, user_uuid): + uid = UUID(str(user_uuid)) + page = int(request.query_params.get("page", 0)) + page_size = int(request.query_params.get("page_size", 20)) + + conv_svc = ConversationService() + read_svc = ReadStateService() + msg_svc = MessageService() + + conversations = conv_svc.list_for_user(uid, page, page_size) + + result = [] + for conv in conversations: + participants = list( + ConversationParticipant.objects.filter(conversation_id=conv.uuid) + .order_by("joined_at") + .values_list("user_uuid", flat=True) + ) + + try: + has_unread = read_svc.has_unread(conv.uuid, uid) + except Exception: + has_unread = False + + try: + last_msgs = msg_svc.list_messages(conv.uuid, page=0, per_page=1) + last_message = last_msgs[0] if last_msgs else None + except Exception: + last_message = None + + result.append( + { + "uuid": conv.uuid, + "type": conv.type, + "status": conv.status, + "closed_by_uuid": conv.closed_by_uuid, + "created_at": conv.created_at, + "participants": participants, + "has_unread": has_unread, + "last_message": last_message, + } + ) + + return Response(ConversationListSerializer(result, many=True).data) + + +class ConversationCloseView(APIView): + authentication_classes = [] + permission_classes = [] + + @extend_schema( + request=ConversationUserActionSerializer, + responses={200: ConversationSerializer}, + ) + def post(self, request, chat_uuid): + serializer = ConversationUserActionSerializer(data=request.data) + serializer.is_valid(raise_exception=True) + + conversation = ConversationService().close( + UUID(str(chat_uuid)), serializer.validated_data["user_uuid"] + ) + + return Response(ConversationSerializer(conversation).data) + + +class ConversationReopenView(APIView): + authentication_classes = [] + permission_classes = [] + + @extend_schema( + request=ConversationUserActionSerializer, + responses={200: ConversationSerializer, 403: None}, + ) + def post(self, request, chat_uuid): + serializer = ConversationUserActionSerializer(data=request.data) + serializer.is_valid(raise_exception=True) + + try: + conversation = ConversationService().reopen( + UUID(str(chat_uuid)), serializer.validated_data["user_uuid"] + ) + except PermissionError as exc: + return Response({"detail": str(exc)}, status=status.HTTP_403_FORBIDDEN) + + return Response(ConversationSerializer(conversation).data) diff --git a/apps/chat/views/messages.py b/apps/chat/views/messages.py new file mode 100644 index 0000000..4ea3da3 --- /dev/null +++ b/apps/chat/views/messages.py @@ -0,0 +1,65 @@ +from uuid import UUID + +from drf_spectacular.utils import OpenApiParameter, extend_schema +from rest_framework import status +from rest_framework.response import Response +from rest_framework.views import APIView + +from apps.chat.exceptions import ConversationClosedError +from apps.chat.serializers.messages import MessageSerializer, SendMessageSerializer +from apps.chat.services.message import MessageService + + +class MessageView(APIView): + authentication_classes = [] + permission_classes = [] + + @extend_schema( + request=SendMessageSerializer, + responses={201: MessageSerializer, 403: None}, + ) + def post(self, request, chat_uuid): + serializer = SendMessageSerializer(data=request.data) + serializer.is_valid(raise_exception=True) + data = serializer.validated_data + + try: + result = MessageService().send( + conversation_uuid=UUID(str(chat_uuid)), + sender_uuid=data["sender_uuid"], + message_type=data["message_type"], + text=data.get("text"), + file=data.get("file"), + ) + except ConversationClosedError as exc: + return Response({"detail": str(exc)}, status=status.HTTP_403_FORBIDDEN) + + return Response(MessageSerializer(result).data, status=status.HTTP_201_CREATED) + + @extend_schema( + parameters=[ + OpenApiParameter("page", int, description="0-based page number", default=0), + OpenApiParameter("per_page", int, description="Results per page", default=20), + OpenApiParameter( + "since", + int, + required=False, + description="Return only posts after this Unix timestamp (milliseconds)", + ), + ], + responses={200: MessageSerializer(many=True)}, + ) + def get(self, request, chat_uuid): + page = int(request.query_params.get("page", 0)) + per_page = int(request.query_params.get("per_page", 20)) + since_raw = request.query_params.get("since") + since = int(since_raw) if since_raw else None + + messages = MessageService().list_messages( + conversation_uuid=UUID(str(chat_uuid)), + page=page, + per_page=per_page, + since=since, + ) + + return Response(MessageSerializer(messages, many=True).data) diff --git a/apps/chat/views/read_state.py b/apps/chat/views/read_state.py new file mode 100644 index 0000000..b477941 --- /dev/null +++ b/apps/chat/views/read_state.py @@ -0,0 +1,70 @@ +import json +from uuid import UUID + +from django.conf import settings +from django_redis import get_redis_connection +from drf_spectacular.utils import extend_schema +from rest_framework.response import Response +from rest_framework.views import APIView + +from apps.chat.serializers.messages import MessageSerializer +from apps.chat.serializers.read_state import MarkReadSerializer, ReadStateSerializer +from apps.chat.services.read_state import ReadStateService + + +class ReadStateView(APIView): + authentication_classes = [] + permission_classes = [] + + @extend_schema( + request=MarkReadSerializer, + responses={200: ReadStateSerializer}, + ) + def post(self, request, chat_uuid): + serializer = MarkReadSerializer(data=request.data) + serializer.is_valid(raise_exception=True) + data = serializer.validated_data + + read_state = ReadStateService().mark_read( + conversation_uuid=UUID(str(chat_uuid)), + user_uuid=data["user_uuid"], + post_id=data["post_id"], + ) + + return Response( + ReadStateSerializer( + { + "conversation_uuid": read_state.conversation_id, + "user_uuid": read_state.user_uuid, + "has_unread": False, + } + ).data + ) + + +class ChatEventsView(APIView): + authentication_classes = [] + permission_classes = [] + + @extend_schema( + responses={200: MessageSerializer(many=True)}, + description=( + "Long-poll endpoint. Blocks up to CHAT_LONG_POLL_TIMEOUT_SECONDS " + "waiting for a new message event. Returns immediately when an event " + "arrives; returns an empty list on timeout." + ), + ) + def get(self, request, chat_uuid): + timeout = getattr(settings, "CHAT_LONG_POLL_TIMEOUT_SECONDS", 25) + + redis = get_redis_connection("default") + key = f"chat_events:{chat_uuid}" + + result = redis.blpop(key, timeout=timeout) + + if result is None: + return Response([]) + + _, raw = result + event = json.loads(raw) + return Response([event]) diff --git a/apps/chat/websocket/__init__.py b/apps/chat/websocket/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/chat/websocket/consumers.py b/apps/chat/websocket/consumers.py new file mode 100644 index 0000000..9298984 --- /dev/null +++ b/apps/chat/websocket/consumers.py @@ -0,0 +1,23 @@ +import json + +from channels.generic.websocket import AsyncWebsocketConsumer + + +class ChatConsumer(AsyncWebsocketConsumer): + async def connect(self): + self.chat_uuid = self.scope["url_route"]["kwargs"]["chat_uuid"] + self.group_name = f"chat_{self.chat_uuid}" + await self.channel_layer.group_add(self.group_name, self.channel_name) + await self.accept() + + async def disconnect(self, close_code): + await self.channel_layer.group_discard(self.group_name, self.channel_name) + + async def receive(self, text_data=None, bytes_data=None): + # Frontend receives only; messages are sent via HTTP API + pass + + async def chat_message(self, event): + """Handle chat.message events forwarded from the channel layer.""" + payload = {k: v for k, v in event.items() if k != "type"} + await self.send(text_data=json.dumps(payload)) diff --git a/apps/chat/websocket/routing.py b/apps/chat/websocket/routing.py new file mode 100644 index 0000000..3bd2b88 --- /dev/null +++ b/apps/chat/websocket/routing.py @@ -0,0 +1,7 @@ +from django.urls import path + +from apps.chat.websocket.consumers import ChatConsumer + +websocket_urlpatterns = [ + path("ws/chat//", ChatConsumer.as_asgi()), +] diff --git a/conftest.py b/conftest.py new file mode 100644 index 0000000..a1dd7ea --- /dev/null +++ b/conftest.py @@ -0,0 +1,22 @@ +import pytest + + +@pytest.fixture(scope="session") +def django_db_setup(django_test_environment, django_db_blocker): + """ + Override pytest-django's default DB setup to: + - Use the local unix-socket superuser instead of the .env credentials + - Remove psycopg3 connection pooling (pool=True conflicts with test + transaction wrapping and causes PoolTimeout during test-DB creation) + """ + from django.conf import settings + from django.test.utils import setup_databases + + db = settings.DATABASES["default"] + db["USER"] = "hashdal" + db["PASSWORD"] = "" + db["HOST"] = "" + db["OPTIONS"] = {} # clear pool=True; psycopg3 pool conflicts with pytest-django + + with django_db_blocker.unblock(): + setup_databases(verbosity=0, interactive=False) diff --git a/env.sample b/env.sample index 3a63c12..1dd3033 100644 --- a/env.sample +++ b/env.sample @@ -13,4 +13,14 @@ REDIS_BASE_URL=redis://127.0.0.1:6379/1 MATTERMOST_TOKEN=3fktq5yzw7yp7m4ty7d3nx8gph MATTERMOST_TOKEN_ID=h5h7unnmeigxdjk4yhtywxjkow -MATTERMOST_BASE_PUBLICE_URL=https://chat.addwin.ir \ No newline at end of file +MATTERMOST_BASE_PUBLICE_URL=https://chat.addwin.ir +MATTERMOST_BASE_PUBLIC_URL=https://chat.addwin.ir +MATTERMOST_TEAM_ID= +MATTERMOST_SERVICE_USERNAME=chat-service +MATTERMOST_SERVICE_EMAIL=chat-service@local.invalid +CHAT_LONG_POLL_TIMEOUT_SECONDS=25 +CHAT_LONG_POLL_INTERVAL_SECONDS=1 +MINIO_ENDPOINT=localhost:9000 +MINIO_ACCESS_KEY=minioadmin +MINIO_SECRET_KEY=minioadmin +MINIO_BUCKET_CHAT=chat diff --git a/main/asgi.py b/main/asgi.py index 04b4e6d..d04d7e7 100644 --- a/main/asgi.py +++ b/main/asgi.py @@ -1,16 +1,27 @@ """ -ASGI config for main project. +ASGI config for chat project. It exposes the ASGI callable as a module-level variable named ``application``. - -For more information on this file, see -https://docs.djangoproject.com/en/5.1/howto/deployment/asgi/ """ import os -from django.core.asgi import get_asgi_application - os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'main.settings') -application = get_asgi_application() +from django.core.asgi import get_asgi_application + +django_asgi_app = get_asgi_application() + +from channels.auth import AuthMiddlewareStack +from channels.routing import ProtocolTypeRouter, URLRouter + +from apps.chat.websocket.routing import websocket_urlpatterns + +application = ProtocolTypeRouter( + { + "http": django_asgi_app, + "websocket": AuthMiddlewareStack( + URLRouter(websocket_urlpatterns) + ), + } +) diff --git a/main/settings.py b/main/settings.py index 177a9c5..47446b4 100644 --- a/main/settings.py +++ b/main/settings.py @@ -56,6 +56,7 @@ INSTALLED_APPS = [ 'apps.core', 'apps.users', 'apps.gooyal_oauth2', + 'apps.chat', ] MIDDLEWARE = [ @@ -135,6 +136,7 @@ AUTHENTICATION_BACKENDS = ( 'django.contrib.auth.backends.ModelBackend', ) WSGI_APPLICATION = 'main.wsgi.application' +ASGI_APPLICATION = 'main.asgi.application' # Database # https://docs.djangoproject.com/en/5.1/ref/settings/#databases @@ -207,6 +209,11 @@ if DEBUG == True: else: STATIC_ROOT = BASE_DIR / 'static' +MEDIA_URL = '/media/' +MEDIA_ROOT = BASE_DIR / 'media' + +DEFAULT_AUTO_FIELD = 'django.db.models.BigAutoField' + # Default primary key field type # https://docs.djangoproject.com/en/5.0/ref/settings/#default-auto-field @@ -216,9 +223,12 @@ SPECTACULAR_SETTINGS = { 'DESCRIPTION': 'service api reference', 'VERSION': '1.0.0', 'SERVE_INCLUDE_SCHEMA': False, - # OTHER SETTINGS - "PARSER_WHITELIST": ["rest_framework.parsers.JSONParser"], - + 'SERVE_PERMISSIONS': ['rest_framework.permissions.AllowAny'], + 'SECURITY': [], + "PARSER_WHITELIST": [ + "rest_framework.parsers.JSONParser", + "rest_framework.parsers.MultiPartParser", + ], } AUTH_USER_MODEL = 'users.User' @@ -247,9 +257,36 @@ CACHES = { } } +CHANNEL_LAYERS = { + "default": { + "BACKEND": "channels_redis.core.RedisChannelLayer", + "CONFIG": { + "hosts": [config('REDIS_BASE_URL')], + }, + } +} + from main.other_settings.logging import get_logging_setting LOGGING = get_logging_setting() NOTIFICATIONS_BASE_PUBLIC_URL = config('NOTIFICATIONS_BASE_PUBLIC_URL', default=None, cast=str) +MATTERMOST_URL = config('MATTERMOST_URL', default=None) MATTERMOST_TOKEN = config('MATTERMOST_TOKEN', default=None, cast=str) +MATTERMOST_TEAM_ID = config('MATTERMOST_TEAM_ID', default=None, cast=str) +# Legacy aliases kept for backward compatibility +MATTERMOST_BASE_URL = config( + 'MATTERMOST_BASE_PUBLIC_URL', + default=config('MATTERMOST_BASE_PUBLICE_URL', default=None, cast=str), + cast=str, +) +MATTERMOST_SERVICE_USERNAME = config('MATTERMOST_SERVICE_USERNAME', default='chat-service', cast=str) +MATTERMOST_SERVICE_EMAIL = config('MATTERMOST_SERVICE_EMAIL', default='chat-service@local.invalid', cast=str) +CHAT_LONG_POLL_TIMEOUT_SECONDS = config('CHAT_LONG_POLL_TIMEOUT_SECONDS', default=25, cast=int) +CHAT_LONG_POLL_INTERVAL_SECONDS = config('CHAT_LONG_POLL_INTERVAL_SECONDS', default=1, cast=int) +MINIO_ENDPOINT = config('MINIO_ENDPOINT', default=None) +MINIO_ACCESS_KEY = config('MINIO_ACCESS_KEY', default=None) +MINIO_SECRET_KEY = config('MINIO_SECRET_KEY', default=None) +MINIO_BUCKET_CHAT = config('MINIO_BUCKET_CHAT', default='chat') +GDAL_LIBRARY_PATH = config('GDAL_LIBRARY_PATH', default=None) +GEOS_LIBRARY_PATH = config('GEOS_LIBRARY_PATH', default=None) diff --git a/main/urls.py b/main/urls.py index 4cdddc4..887d94e 100644 --- a/main/urls.py +++ b/main/urls.py @@ -36,6 +36,7 @@ urlpatterns = [ path('admin/', admin.site.urls), # path('', include('apps.pg.urls')), # path('users/', include('apps.users.urls')), + path('', include('apps.chat.urls')), path('oauth2/', include('oauth2_provider.urls', namespace='oauth2_provider')), ] diff --git a/pytest.ini b/pytest.ini new file mode 100644 index 0000000..71ddc8e --- /dev/null +++ b/pytest.ini @@ -0,0 +1,2 @@ +[pytest] +DJANGO_SETTINGS_MODULE = main.settings diff --git a/requirements.in b/requirements.in index 1d1057a..109e9de 100644 --- a/requirements.in +++ b/requirements.in @@ -20,4 +20,9 @@ django-redis xlsxwriter xlrd django_admin_logs -mattermostdriver \ No newline at end of file +mattermostdriver +channels +channels-redis +minio +pytest +pytest-django \ No newline at end of file diff --git a/requirements.txt b/requirements.txt index eed64ab..a4ea010 100644 --- a/requirements.txt +++ b/requirements.txt @@ -8,6 +8,10 @@ amqp==5.3.1 # via kombu anyio==4.13.0 # via httpx +argon2-cffi==25.1.0 + # via minio +argon2-cffi-bindings==25.1.0 + # via argon2-cffi asgiref==3.11.1 # via # django @@ -26,7 +30,15 @@ certifi==2026.4.22 # httpx # requests cffi==2.0.0 - # via cryptography + # via + # argon2-cffi-bindings + # cryptography +channels==4.3.2 + # via + # -r requirements.in + # channels-redis +channels-redis==4.3.0 + # via -r requirements.in charset-normalizer==3.4.7 # via requests click==8.3.3 @@ -99,6 +111,8 @@ idna==3.13 # requests inflection==0.5.1 # via drf-spectacular +iniconfig==2.3.0 + # via pytest jalali-core==1.0.0 # via jdatetime jdatetime==5.2.0 @@ -113,6 +127,10 @@ kombu==5.6.2 # via celery mattermostdriver==7.3.2 # via -r requirements.in +minio==7.2.20 + # via -r requirements.in +msgpack==1.1.2 + # via channels-redis oauthlib==3.3.1 # via django-oauth-toolkit packaging==26.2 @@ -121,6 +139,8 @@ packaging==26.2 # kombu pillow==12.2.0 # via -r requirements.in +pluggy==1.6.0 + # via pytest prompt-toolkit==3.0.52 # via click-repl psycopg[binary,pool]==3.3.4 @@ -131,6 +151,16 @@ psycopg-pool==3.3.1 # via psycopg pycparser==3.0 # via cffi +pycryptodome==3.23.0 + # via minio +pygments==2.20.0 + # via pytest +pytest==9.0.3 + # via + # -r requirements.in + # pytest-django +pytest-django==4.12.0 + # via -r requirements.in python-dateutil==2.9.0.post0 # via celery python-decouple==3.8