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