Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2c3d691f27 | ||
|
|
0dd06ce89d |
+22
-8
@@ -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);
|
||||||
@@ -67,7 +80,8 @@ 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.
Binary file not shown.
Binary file not shown.
Binary file not shown.
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")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user