feature/refactor_uuid #1

Merged
Ghasemi merged 9 commits from feature/refactor_uuid into master 2026-07-19 05:49:06 -04:00
58 changed files with 1896 additions and 22 deletions

41
.env.example Normal file
View file

@ -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=

5
.gitignore vendored
View file

@ -5,3 +5,8 @@ media
/delme.py
/log/accounts.log
/log/errors.log
**/__pycache__
**/*.pyc
**/*.pyo
.claude
/log

View file

@ -3,3 +3,4 @@ from django.apps import AppConfig
class ChatConfig(AppConfig):
name = 'apps.chat'
default_auto_field = 'django.db.models.BigAutoField'

View file

0
apps/chat/events/base.py Normal file
View file

11
apps/chat/events/event.py Normal file
View file

@ -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

View file

View file

@ -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)

View file

@ -0,0 +1,7 @@
from apps.chat.events.event import MessageSentEvent
class PushPublisher:
# Future: FCM / APNS integration
def publish(self, event: MessageSentEvent) -> None:
pass

View file

@ -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,
},
)

2
apps/chat/exceptions.py Normal file
View file

@ -0,0 +1,2 @@
class ConversationClosedError(Exception):
"""Raised when attempting to send a message to a closed conversation."""

View file

View file

@ -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

View file

@ -0,0 +1,2 @@
class MattermostError(Exception):
"""Raised when any Mattermost API call fails."""

View file

@ -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')],
},
),
]

View file

@ -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'),
),
]

View file

@ -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'),
),
]

View file

@ -1,3 +0,0 @@
from django.db import models
# Create your models here.

View file

@ -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",
]

View file

@ -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)

View file

@ -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}"

View file

@ -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}"

View file

@ -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}"

View file

View file

@ -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)

View file

@ -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)

View file

@ -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()

View file

View file

@ -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

View file

@ -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

View file

@ -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,
}

View file

@ -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

View file

@ -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},
)

View file

@ -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}"

View file

@ -1,3 +0,0 @@
from django.test import TestCase
# Create your tests here.

View file

View file

@ -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")

View file

@ -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

View file

@ -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

View file

@ -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

20
apps/chat/urls.py Normal file
View file

@ -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/<uuid:user_uuid>/chats/", UserConversationListView.as_view(), name="user-chat-list"),
path("api/chats/<uuid:chat_uuid>/messages/", MessageView.as_view(), name="chat-messages"),
path("api/chats/<uuid:chat_uuid>/read/", ReadStateView.as_view(), name="chat-read"),
path("api/chats/<uuid:chat_uuid>/events/", ChatEventsView.as_view(), name="chat-events"),
path("api/chats/<uuid:chat_uuid>/close/", ConversationCloseView.as_view(), name="chat-close"),
path("api/chats/<uuid:chat_uuid>/reopen/", ConversationReopenView.as_view(), name="chat-reopen"),
]

View file

@ -1,3 +0,0 @@
from django.shortcuts import render
# Create your views here.

View file

View file

@ -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)

View file

@ -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)

View file

@ -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])

View file

View file

@ -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))

View file

@ -0,0 +1,7 @@
from django.urls import path
from apps.chat.websocket.consumers import ChatConsumer
websocket_urlpatterns = [
path("ws/chat/<uuid:chat_uuid>/", ChatConsumer.as_asgi()),
]

22
conftest.py Normal file
View file

@ -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)

View file

@ -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
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

View file

@ -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)
),
}
)

View file

@ -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)

View file

@ -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')),
]

2
pytest.ini Normal file
View file

@ -0,0 +1,2 @@
[pytest]
DJANGO_SETTINGS_MODULE = main.settings

View file

@ -20,4 +20,9 @@ django-redis
xlsxwriter
xlrd
django_admin_logs
mattermostdriver
mattermostdriver
channels
channels-redis
minio
pytest
pytest-django

View file

@ -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