2 Commits
13 changed files with 64 additions and 34 deletions
@@ -1,13 +1,14 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using MediatR;
using System.Linq;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Modules.Conversations.Infrastructure.SignalR;
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Storage;
using MediatR;
using Microsoft.AspNetCore.SignalR;
namespace Knot.Modules.Conversations.Application.Chats.LeaveOrDelete;
@@ -19,17 +20,21 @@ internal sealed class LeaveOrDeleteChatCommandHandler : ICommandHandler<LeaveOrD
private readonly IMessageRepository _messageRepository;
private readonly IFileStorageService _fileStorage;
private readonly IChatsUnitOfWork _uow;
private readonly IHubContext<ChatHub> _hubContext;
public LeaveOrDeleteChatCommandHandler(
IChatRepository chatRepository,
IMessageRepository messageRepository,
IFileStorageService fileStorage,
IChatsUnitOfWork uow)
IChatsUnitOfWork uow,
IHubContext<ChatHub> hubContext)
{
_chatRepository = chatRepository;
_messageRepository = messageRepository;
_fileStorage = fileStorage;
_uow = uow;
_hubContext = hubContext;
}
public async Task<Result<SuccessResponse>> Handle(LeaveOrDeleteChatCommand request, CancellationToken cancellationToken)
@@ -58,6 +63,14 @@ internal sealed class LeaveOrDeleteChatCommandHandler : ICommandHandler<LeaveOrD
// DELETE ALL MESSAGES AND FILES FIRST
await DeleteChatMediaAndMessagesAsync(chat.Id, cancellationToken);
_chatRepository.Remove(chat);
// Notify all remaining members that the chat was deleted
foreach (var member in chat.Members)
{
await _hubContext.Clients.User(member.UserId.ToString())
.SendAsync("chat_deleted", chat.Id.ToString(), cancellationToken);
}
}
await _uow.SaveChangesAsync(cancellationToken);
@@ -68,6 +81,7 @@ internal sealed class LeaveOrDeleteChatCommandHandler : ICommandHandler<LeaveOrD
private async Task DeleteChatMediaAndMessagesAsync(Guid chatId, CancellationToken ct)
{
try
{
// Get all messages directly from Mongo (not paged)
var messages = await _messageRepository.GetChatMessagesAsync(chatId, int.MaxValue, 0, ct);
Binary file not shown.
Binary file not shown.
@@ -7,6 +7,7 @@ sealed class ChatEvent {
data class NewMessage(val message: MessageDto) : ChatEvent()
data class MessageEdited(val messageId: String, val chatId: String, val content: String) : ChatEvent()
data class MessageDeleted(val messageId: String, val chatId: String) : ChatEvent()
data class ChatDeleted(val chatId: String) : ChatEvent()
data class MessagesRead(val chatId: String, val userId: String, val lastReadSequenceId: Int) : ChatEvent()
data class UserTyping(val chatId: String, val userId: String) : ChatEvent()
data class UserStoppedTyping(val chatId: String, val userId: String) : ChatEvent()
@@ -136,6 +136,11 @@ class ChatHubClient @Inject constructor() {
_events.tryEmit(ChatEvent.MessageDeleted(messageId, chatId))
}, String::class.java, String::class.java)
conn.on("chat_deleted", { chatId: String ->
Log.d("ChatHubClient", ">>> chat_deleted event: $chatId")
_events.tryEmit(ChatEvent.ChatDeleted(chatId))
}, String::class.java)
conn.on("messages_read", { data: MessagesReadEvent ->
Log.d("ChatHubClient", ">>> messages_read event: ${data.effectiveChatId}")
_events.tryEmit(ChatEvent.MessagesRead(
@@ -83,6 +83,18 @@ class ChatRepositoryImpl @Inject constructor(
}
}
android.util.Log.d(TAG, "Cached ${entities.size} chats with messages")
// Проверяем удалённые чаты - если чата нет в списке от сервера, удаляем из БД
val serverChatIds = chats.map { it.id }.toSet()
val localChats = chatDao.getAllChats()
localChats.forEach { localChat ->
if (localChat.id !in serverChatIds) {
android.util.Log.d(TAG, "Chat ${localChat.id} was deleted on server, removing from local DB")
chatDao.deleteChat(localChat.id)
messageDao.clearChat(localChat.id)
}
}
chats
} catch (e: Exception) {
android.util.Log.d(TAG, "Network load failed, using cache")
@@ -8,6 +8,7 @@ import chats.data.remote.signalr.MessagesReadEvent
import chats.data.remote.signalr.ReactionEvent
import chats.data.repository.toEntity
import com.google.gson.Gson
import core.database.data.ChatDao
import core.database.data.MessageDao
import core.database.data.MessageEntity
import core.database.data.SyncStatus
@@ -18,22 +19,11 @@ import kotlinx.coroutines.flow.collectLatest
import javax.inject.Inject
import javax.inject.Singleton
/**
* Обработчик SignalR событий для обновления локального кэша
*
* Обрабатывает события:
* - new_message: новое сообщение в чате
* - message_edited: сообщение отредактировано
* - message_deleted: сообщение удалено
* - messages_read: сообщения прочитаны
* - reaction_added/removed: реакция добавлена/удалена
*
* Все изменения сразу записываются в Room, UI обновляется через Flow
*/
@Singleton
class MessageSignalRHandler @Inject constructor(
private val hubClient: ChatHubClient,
private val dao: MessageDao,
private val messageDao: MessageDao,
private val chatDao: ChatDao,
private val serverConfig: ServerConfig,
private val tokenManager: TokenManager
) {
@@ -78,6 +68,7 @@ class MessageSignalRHandler @Inject constructor(
is ChatEvent.NewMessage -> handleNewMessage(event.message)
is ChatEvent.MessageEdited -> handleMessageEdited(event.messageId, event.chatId, event.content)
is ChatEvent.MessageDeleted -> handleMessageDeleted(event.messageId, event.chatId)
is ChatEvent.ChatDeleted -> handleChatDeleted(event.chatId)
is ChatEvent.MessagesRead -> handleMessagesRead(event.chatId, event.lastReadSequenceId)
is ChatEvent.ReactionUpdated -> handleReactionUpdated(
event.messageId,
@@ -91,6 +82,12 @@ class MessageSignalRHandler @Inject constructor(
}
}
private suspend fun handleChatDeleted(chatId: String) {
Log.d(TAG, "Chat deleted: $chatId, removing from local database")
chatDao.deleteChat(chatId)
Log.d(TAG, "Successfully removed chat $chatId from local database")
}
private suspend fun handleNewMessage(dto: MessageDto) {
Log.d(TAG, "New message received: ${dto.id} in chat ${dto.chatId}")
@@ -99,7 +96,7 @@ class MessageSignalRHandler @Inject constructor(
syncStatus = SyncStatus.SYNCED
)
val existing = dao.getMessageById(dto.id)
val existing = messageDao.getMessageById(dto.id)
if (existing != null && existing.syncStatus == SyncStatus.SYNCING) {
val syncedEntity = entity.copy(
syncStatus = SyncStatus.SYNCED,
@@ -107,10 +104,10 @@ class MessageSignalRHandler @Inject constructor(
isEditedLocally = false,
editedContent = null
)
dao.insertMessage(syncedEntity)
messageDao.insertMessage(syncedEntity)
Log.d(TAG, "Merged local message with server response: ${dto.id}")
} else {
dao.insertMessage(entity)
messageDao.insertMessage(entity)
Log.d(TAG, "Saved new message: ${dto.id}")
}
}
@@ -118,7 +115,7 @@ class MessageSignalRHandler @Inject constructor(
private suspend fun handleMessageEdited(messageId: String, chatId: String, content: String) {
Log.d(TAG, "Message edited: $messageId")
val existing = dao.getMessageById(messageId)
val existing = messageDao.getMessageById(messageId)
if (existing != null) {
if (existing.isEditedLocally) {
Log.d(TAG, "Message is being edited locally, skipping server update: $messageId")
@@ -129,7 +126,7 @@ class MessageSignalRHandler @Inject constructor(
content = content,
lastUpdated = System.currentTimeMillis()
)
dao.insertMessage(updated)
messageDao.insertMessage(updated)
Log.d(TAG, "Updated edited message: $messageId")
} else {
Log.w(TAG, "Edited message not found in cache: $messageId")
@@ -139,12 +136,12 @@ class MessageSignalRHandler @Inject constructor(
private suspend fun handleMessageDeleted(messageId: String, chatId: String) {
Log.d(TAG, "Message deleted: $messageId")
val existing = dao.getMessageById(messageId)
val existing = messageDao.getMessageById(messageId)
if (existing?.isDeletedLocally == true) {
dao.deleteMessage(messageId)
messageDao.deleteMessage(messageId)
Log.d(TAG, "Completed local deletion: $messageId")
} else {
dao.deleteMessage(messageId)
messageDao.deleteMessage(messageId)
Log.d(TAG, "Removed deleted message from cache: $messageId")
}
}
@@ -156,13 +153,13 @@ class MessageSignalRHandler @Inject constructor(
}
Log.d(TAG, "Messages read in chat $chatId up to sequenceId $lastReadSequenceId")
dao.markMessagesAsRead(chatId, lastReadSequenceId)
messageDao.markMessagesAsRead(chatId, lastReadSequenceId)
}
private suspend fun handleReactionUpdated(messageId: String, emoji: String, isRemoved: Boolean) {
Log.d(TAG, "Reaction ${if (isRemoved) "removed" else "added"}: $emoji on message $messageId")
val existing = dao.getMessageById(messageId)
val existing = messageDao.getMessageById(messageId)
if (existing != null) {
val reactions = gson.fromJson<Map<String, Int>>(existing.reactionsJson, Map::class.java) ?: emptyMap()
val updatedReactions = reactions.toMutableMap()
@@ -182,7 +179,7 @@ class MessageSignalRHandler @Inject constructor(
reactionsJson = gson.toJson(updatedReactions),
lastUpdated = System.currentTimeMillis()
)
dao.insertMessage(updated)
messageDao.insertMessage(updated)
Log.d(TAG, "Updated reactions for message: $messageId")
}
}
+3 -2
View File
@@ -33,11 +33,12 @@ object ChatModule {
@Singleton
fun provideMessageSignalRHandler(
hubClient: ChatHubClient,
dao: MessageDao,
messageDao: MessageDao,
chatDao: ChatDao,
serverConfig: ServerConfig,
tokenManager: TokenManager
): MessageSignalRHandler {
return MessageSignalRHandler(hubClient, dao, serverConfig, tokenManager)
return MessageSignalRHandler(hubClient, messageDao, chatDao, serverConfig, tokenManager)
}
@Provides