Merge pull request 'feature/refactor_uuid' (#1) from feature/refactor_uuid into master
Reviewed-on: #1
This commit is contained in:
commit
0802bc93d6
58 changed files with 1896 additions and 22 deletions
41
.env.example
Normal file
41
.env.example
Normal 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
5
.gitignore
vendored
|
|
@ -5,3 +5,8 @@ media
|
|||
/delme.py
|
||||
/log/accounts.log
|
||||
/log/errors.log
|
||||
**/__pycache__
|
||||
**/*.pyc
|
||||
**/*.pyo
|
||||
.claude
|
||||
/log
|
||||
|
|
@ -3,3 +3,4 @@ from django.apps import AppConfig
|
|||
|
||||
class ChatConfig(AppConfig):
|
||||
name = 'apps.chat'
|
||||
default_auto_field = 'django.db.models.BigAutoField'
|
||||
|
|
|
|||
0
apps/chat/events/__init__.py
Normal file
0
apps/chat/events/__init__.py
Normal file
0
apps/chat/events/base.py
Normal file
0
apps/chat/events/base.py
Normal file
11
apps/chat/events/event.py
Normal file
11
apps/chat/events/event.py
Normal 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
|
||||
0
apps/chat/events/publishers/__init__.py
Normal file
0
apps/chat/events/publishers/__init__.py
Normal file
23
apps/chat/events/publishers/longpoll.py
Normal file
23
apps/chat/events/publishers/longpoll.py
Normal 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)
|
||||
7
apps/chat/events/publishers/push.py
Normal file
7
apps/chat/events/publishers/push.py
Normal 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
|
||||
20
apps/chat/events/publishers/websocket.py
Normal file
20
apps/chat/events/publishers/websocket.py
Normal 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
2
apps/chat/exceptions.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
class ConversationClosedError(Exception):
|
||||
"""Raised when attempting to send a message to a closed conversation."""
|
||||
0
apps/chat/integrations/__init__.py
Normal file
0
apps/chat/integrations/__init__.py
Normal file
0
apps/chat/integrations/mattermost/__init__.py
Normal file
0
apps/chat/integrations/mattermost/__init__.py
Normal file
218
apps/chat/integrations/mattermost/client.py
Normal file
218
apps/chat/integrations/mattermost/client.py
Normal 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
|
||||
2
apps/chat/integrations/mattermost/exceptions.py
Normal file
2
apps/chat/integrations/mattermost/exceptions.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
class MattermostError(Exception):
|
||||
"""Raised when any Mattermost API call fails."""
|
||||
68
apps/chat/migrations/0001_initial.py
Normal file
68
apps/chat/migrations/0001_initial.py
Normal 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')],
|
||||
},
|
||||
),
|
||||
]
|
||||
|
|
@ -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'),
|
||||
),
|
||||
]
|
||||
|
|
@ -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'),
|
||||
),
|
||||
]
|
||||
|
|
@ -1,3 +0,0 @@
|
|||
from django.db import models
|
||||
|
||||
# Create your models here.
|
||||
13
apps/chat/models/__init__.py
Normal file
13
apps/chat/models/__init__.py
Normal 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",
|
||||
]
|
||||
43
apps/chat/models/conversation.py
Normal file
43
apps/chat/models/conversation.py
Normal 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)
|
||||
18
apps/chat/models/mapping.py
Normal file
18
apps/chat/models/mapping.py
Normal 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}"
|
||||
31
apps/chat/models/participant.py
Normal file
31
apps/chat/models/participant.py
Normal 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}"
|
||||
32
apps/chat/models/read_state.py
Normal file
32
apps/chat/models/read_state.py
Normal 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}"
|
||||
0
apps/chat/serializers/__init__.py
Normal file
0
apps/chat/serializers/__init__.py
Normal file
40
apps/chat/serializers/conversations.py
Normal file
40
apps/chat/serializers/conversations.py
Normal 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)
|
||||
17
apps/chat/serializers/messages.py
Normal file
17
apps/chat/serializers/messages.py
Normal 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)
|
||||
12
apps/chat/serializers/read_state.py
Normal file
12
apps/chat/serializers/read_state.py
Normal 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()
|
||||
0
apps/chat/services/__init__.py
Normal file
0
apps/chat/services/__init__.py
Normal file
7
apps/chat/services/account.py
Normal file
7
apps/chat/services/account.py
Normal 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
|
||||
91
apps/chat/services/conversation.py
Normal file
91
apps/chat/services/conversation.py
Normal 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
|
||||
128
apps/chat/services/message.py
Normal file
128
apps/chat/services/message.py
Normal 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,
|
||||
}
|
||||
48
apps/chat/services/read_state.py
Normal file
48
apps/chat/services/read_state.py
Normal 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
|
||||
12
apps/chat/services/realtime.py
Normal file
12
apps/chat/services/realtime.py
Normal 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},
|
||||
)
|
||||
42
apps/chat/services/storage.py
Normal file
42
apps/chat/services/storage.py
Normal 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}"
|
||||
|
|
@ -1,3 +0,0 @@
|
|||
from django.test import TestCase
|
||||
|
||||
# Create your tests here.
|
||||
0
apps/chat/tests/__init__.py
Normal file
0
apps/chat/tests/__init__.py
Normal file
112
apps/chat/tests/test_conversation_close.py
Normal file
112
apps/chat/tests/test_conversation_close.py
Normal 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")
|
||||
73
apps/chat/tests/test_conversations.py
Normal file
73
apps/chat/tests/test_conversations.py
Normal 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
|
||||
116
apps/chat/tests/test_messages.py
Normal file
116
apps/chat/tests/test_messages.py
Normal 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
|
||||
86
apps/chat/tests/test_read_state.py
Normal file
86
apps/chat/tests/test_read_state.py
Normal 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
20
apps/chat/urls.py
Normal 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"),
|
||||
]
|
||||
|
|
@ -1,3 +0,0 @@
|
|||
from django.shortcuts import render
|
||||
|
||||
# Create your views here.
|
||||
0
apps/chat/views/__init__.py
Normal file
0
apps/chat/views/__init__.py
Normal file
138
apps/chat/views/conversations.py
Normal file
138
apps/chat/views/conversations.py
Normal 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)
|
||||
65
apps/chat/views/messages.py
Normal file
65
apps/chat/views/messages.py
Normal 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)
|
||||
70
apps/chat/views/read_state.py
Normal file
70
apps/chat/views/read_state.py
Normal 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])
|
||||
0
apps/chat/websocket/__init__.py
Normal file
0
apps/chat/websocket/__init__.py
Normal file
23
apps/chat/websocket/consumers.py
Normal file
23
apps/chat/websocket/consumers.py
Normal 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))
|
||||
7
apps/chat/websocket/routing.py
Normal file
7
apps/chat/websocket/routing.py
Normal 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
22
conftest.py
Normal 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)
|
||||
12
env.sample
12
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
|
||||
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
|
||||
|
|
|
|||
25
main/asgi.py
25
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)
|
||||
),
|
||||
}
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
2
pytest.ini
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
[pytest]
|
||||
DJANGO_SETTINGS_MODULE = main.settings
|
||||
|
|
@ -20,4 +20,9 @@ django-redis
|
|||
xlsxwriter
|
||||
xlrd
|
||||
django_admin_logs
|
||||
mattermostdriver
|
||||
mattermostdriver
|
||||
channels
|
||||
channels-redis
|
||||
minio
|
||||
pytest
|
||||
pytest-django
|
||||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue