Архитектура
This commit is contained in:
@@ -18,12 +18,14 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
|
||||
private readonly IChatRepository _chatRepository;
|
||||
private readonly IUserDisplayNameProvider _userProvider;
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IMessageReactionRepository _reactionRepository;
|
||||
|
||||
public GetChatByIdQueryHandler(IChatRepository chatRepository, IUserDisplayNameProvider userProvider, IMessageRepository messageRepository)
|
||||
public GetChatByIdQueryHandler(IChatRepository chatRepository, IUserDisplayNameProvider userProvider, IMessageRepository messageRepository, IMessageReactionRepository reactionRepository)
|
||||
{
|
||||
_chatRepository = chatRepository;
|
||||
_userProvider = userProvider;
|
||||
_messageRepository = messageRepository;
|
||||
_reactionRepository = reactionRepository;
|
||||
}
|
||||
|
||||
public async Task<Result<ChatDto?>> Handle(GetChatByIdQuery request, CancellationToken cancellationToken)
|
||||
@@ -46,10 +48,14 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
|
||||
}
|
||||
|
||||
var latestMessage = await _messageRepository.GetLatestChatMessageAsync(chat.Id, cancellationToken);
|
||||
var latestReactions = latestMessage != null
|
||||
? await _reactionRepository.GetReactionsForMessageAsync(latestMessage.Id, cancellationToken)
|
||||
: new List<MessageReaction>();
|
||||
|
||||
if (latestMessage != null)
|
||||
{
|
||||
userIdsToFetch.Add(latestMessage.SenderId);
|
||||
foreach (var reaction in latestMessage.Reactions)
|
||||
foreach (var reaction in latestReactions)
|
||||
{
|
||||
userIdsToFetch.Add(reaction.UserId);
|
||||
}
|
||||
@@ -83,7 +89,7 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
|
||||
usersInfo.TryGetValue(latestMessage.SenderId, out var senderObj);
|
||||
|
||||
var reactionsWithUser = new List<ReactionDto>();
|
||||
foreach (var reaction in latestMessage.Reactions)
|
||||
foreach (var reaction in latestReactions)
|
||||
{
|
||||
usersInfo.TryGetValue(reaction.UserId, out var reactionUser);
|
||||
reactionsWithUser.Add(new ReactionDto(
|
||||
@@ -96,6 +102,11 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
|
||||
));
|
||||
}
|
||||
|
||||
var readByList = chat.Members
|
||||
.Where(m => m.LastReadSequenceId >= latestMessage.SequenceId && m.UserId != latestMessage.SenderId)
|
||||
.Select(m => new ReadByDto(m.UserId))
|
||||
.ToList();
|
||||
|
||||
messagesList.Add(new ChatMessageDto(
|
||||
latestMessage.Id,
|
||||
latestMessage.ChatId,
|
||||
@@ -110,6 +121,7 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
|
||||
latestMessage.IsEdited,
|
||||
latestMessage.IsDeleted,
|
||||
latestMessage.CreatedAt,
|
||||
latestMessage.SequenceId,
|
||||
latestMessage.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList(),
|
||||
senderObj != null ? new MessageSenderDto(
|
||||
senderObj.Id,
|
||||
@@ -118,10 +130,13 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
|
||||
senderObj.Avatar
|
||||
) : new MessageSenderDto(latestMessage.SenderId, "unknown", "Unknown", null),
|
||||
reactionsWithUser,
|
||||
latestMessage.ReadBy.Select(readReceipt => new ReadByDto(readReceipt.UserId)).ToList()
|
||||
readByList
|
||||
));
|
||||
}
|
||||
|
||||
var currentMember = chat.Members.First(m => m.UserId == request.UserId);
|
||||
var unreadCount = (int)Math.Max(0, chat.LastMessageSequenceId - currentMember.LastReadSequenceId);
|
||||
|
||||
var dto = new ChatDto(
|
||||
chat.Id,
|
||||
chat.Type.ToString().ToLowerInvariant(),
|
||||
@@ -131,7 +146,7 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
|
||||
chat.CreatedAt,
|
||||
members,
|
||||
messagesList,
|
||||
0
|
||||
unreadCount
|
||||
);
|
||||
|
||||
return Result.Success<ChatDto?>(dto);
|
||||
|
||||
@@ -19,12 +19,14 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
|
||||
private readonly IChatRepository _chatRepository;
|
||||
private readonly IUserDisplayNameProvider _userProvider;
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IMessageReactionRepository _reactionRepository;
|
||||
|
||||
public GetChatsQueryHandler(IChatRepository chatRepository, IUserDisplayNameProvider userProvider, IMessageRepository messageRepository)
|
||||
public GetChatsQueryHandler(IChatRepository chatRepository, IUserDisplayNameProvider userProvider, IMessageRepository messageRepository, IMessageReactionRepository reactionRepository)
|
||||
{
|
||||
_chatRepository = chatRepository;
|
||||
_userProvider = userProvider;
|
||||
_messageRepository = messageRepository;
|
||||
_reactionRepository = reactionRepository;
|
||||
}
|
||||
|
||||
public async Task<Result<List<ChatDto>>> Handle(GetChatsQuery request, CancellationToken cancellationToken)
|
||||
@@ -35,6 +37,10 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
|
||||
foreach (var chat in userChats)
|
||||
{
|
||||
var latestMessage = await _messageRepository.GetLatestChatMessageAsync(chat.Id, cancellationToken);
|
||||
var latestReactions = latestMessage != null
|
||||
? await _reactionRepository.GetReactionsForMessageAsync(latestMessage.Id, cancellationToken)
|
||||
: new List<MessageReaction>();
|
||||
|
||||
var userIdsToFetch = new HashSet<Guid>();
|
||||
|
||||
foreach (var member in chat.Members)
|
||||
@@ -45,7 +51,7 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
|
||||
if (latestMessage != null)
|
||||
{
|
||||
userIdsToFetch.Add(latestMessage.SenderId);
|
||||
foreach (var r in latestMessage.Reactions)
|
||||
foreach (var r in latestReactions)
|
||||
{
|
||||
userIdsToFetch.Add(r.UserId);
|
||||
}
|
||||
@@ -80,7 +86,7 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
|
||||
usersInfo.TryGetValue(latestMessage.SenderId, out var senderObj);
|
||||
|
||||
var reactionsWithUser = new List<ReactionDto>();
|
||||
foreach (var reaction in latestMessage.Reactions)
|
||||
foreach (var reaction in latestReactions)
|
||||
{
|
||||
usersInfo.TryGetValue(reaction.UserId, out var reactionUser);
|
||||
reactionsWithUser.Add(new ReactionDto(
|
||||
@@ -107,6 +113,7 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
|
||||
latestMessage.IsEdited,
|
||||
latestMessage.IsDeleted,
|
||||
latestMessage.CreatedAt,
|
||||
latestMessage.SequenceId,
|
||||
latestMessage.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList(),
|
||||
senderObj != null ? new MessageSenderDto(
|
||||
senderObj.Id,
|
||||
@@ -115,11 +122,12 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
|
||||
senderObj.Avatar
|
||||
) : new MessageSenderDto(latestMessage.SenderId, "unknown", "Unknown", null),
|
||||
reactionsWithUser,
|
||||
latestMessage.ReadBy.Select(readReceipt => new ReadByDto(readReceipt.UserId)).ToList()
|
||||
chat.Members.Where(m => m.LastReadSequenceId >= latestMessage.SequenceId && m.UserId != latestMessage.SenderId).Select(m => new ReadByDto(m.UserId)).ToList()
|
||||
));
|
||||
}
|
||||
|
||||
var unreadCount = await _messageRepository.GetUnreadCountAsync(chat.Id, request.UserId, cancellationToken);
|
||||
var currentMember = chat.Members.FirstOrDefault(m => m.UserId == request.UserId);
|
||||
var unreadCount = currentMember != null ? (int)Math.Max(0, chat.LastMessageSequenceId - currentMember.LastReadSequenceId) : 0;
|
||||
|
||||
dtos.Add(new ChatDto(
|
||||
chat.Id,
|
||||
|
||||
@@ -17,6 +17,7 @@ public record ChatMessageDto(
|
||||
bool IsEdited,
|
||||
bool IsDeleted,
|
||||
DateTime CreatedAt,
|
||||
long SequenceId,
|
||||
List<MediaDto> Media,
|
||||
MessageSenderDto Sender,
|
||||
List<ReactionDto> Reactions,
|
||||
|
||||
@@ -15,6 +15,7 @@ public record MessageDetailDto(
|
||||
bool IsEdited,
|
||||
bool IsDeleted,
|
||||
DateTime CreatedAt,
|
||||
long SequenceId,
|
||||
Guid? ForwardedFromId,
|
||||
MessageSenderDto? ForwardedFrom,
|
||||
Guid? StoryId,
|
||||
|
||||
@@ -14,6 +14,7 @@ public record SearchMessageDto(
|
||||
bool IsEdited,
|
||||
bool IsDeleted,
|
||||
DateTime CreatedAt,
|
||||
long SequenceId,
|
||||
Guid? ForwardedFromId,
|
||||
MessageSenderDto? ForwardedFrom,
|
||||
Guid? StoryId,
|
||||
|
||||
@@ -19,12 +19,14 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IUserDisplayNameProvider _userProvider;
|
||||
private readonly IChatRepository _chatRepository;
|
||||
private readonly IMessageReactionRepository _reactionRepository;
|
||||
|
||||
public GetMessagesQueryHandler(IMessageRepository messageRepository, IUserDisplayNameProvider userProvider, IChatRepository chatRepository)
|
||||
public GetMessagesQueryHandler(IMessageRepository messageRepository, IUserDisplayNameProvider userProvider, IChatRepository chatRepository, IMessageReactionRepository reactionRepository)
|
||||
{
|
||||
_messageRepository = messageRepository;
|
||||
_userProvider = userProvider;
|
||||
_chatRepository = chatRepository;
|
||||
_reactionRepository = reactionRepository;
|
||||
}
|
||||
|
||||
public async Task<Result<List<MessageDetailDto>>> Handle(GetMessagesQuery request, CancellationToken cancellationToken)
|
||||
@@ -70,6 +72,10 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
|
||||
}
|
||||
|
||||
var senders = await _userProvider.GetUsersInfoAsync(userIdsToFetch, cancellationToken);
|
||||
|
||||
var messageIds = filteredMessages.Select(m => m.Id).ToList();
|
||||
var allReactions = await _reactionRepository.GetReactionsForMessagesAsync(messageIds, cancellationToken);
|
||||
var reactionsByMessage = allReactions.GroupBy(r => r.MessageId).ToDictionary(g => g.Key, g => g.ToList());
|
||||
|
||||
foreach (var message in messages)
|
||||
{
|
||||
@@ -95,7 +101,8 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
|
||||
}
|
||||
|
||||
var reactionsWithUser = new List<MessageReactionDto>();
|
||||
foreach (var reaction in message.Reactions)
|
||||
var messageReactions = reactionsByMessage.TryGetValue(message.Id, out var mr) ? mr : new List<MessageReaction>();
|
||||
foreach (var reaction in messageReactions)
|
||||
{
|
||||
var userObj = senders.TryGetValue(reaction.UserId, out var reactionUser)
|
||||
? new MessageSenderDto(reactionUser.Id, reactionUser.Username, reactionUser.DisplayName, reactionUser.Avatar)
|
||||
@@ -121,6 +128,7 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
|
||||
message.IsEdited,
|
||||
message.IsDeleted,
|
||||
message.CreatedAt,
|
||||
message.SequenceId,
|
||||
message.ForwardedFromId,
|
||||
message.ForwardedFromId.HasValue && senders.TryGetValue(message.ForwardedFromId.Value, out var fwdUser) ? new MessageSenderDto(fwdUser.Id, fwdUser.Username, fwdUser.DisplayName, fwdUser.Avatar) : null,
|
||||
message.StoryId,
|
||||
@@ -128,7 +136,7 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
|
||||
message.StoryMediaType,
|
||||
message.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList(),
|
||||
senders.TryGetValue(message.SenderId, out var senderUser) ? new MessageSenderDto(senderUser.Id, senderUser.Username, senderUser.DisplayName, senderUser.Avatar) : null,
|
||||
message.ReadBy.Select(readReceipt => new ReadByDto(readReceipt.UserId)).ToList(),
|
||||
chat.Members.Where(m => m.LastReadSequenceId >= message.SequenceId && m.UserId != message.SenderId).Select(m => new ReadByDto(m.UserId)).ToList(),
|
||||
reactionsWithUser
|
||||
));
|
||||
}
|
||||
|
||||
@@ -92,6 +92,11 @@ internal sealed class GetSharedMediaQueryHandler : IQueryHandler<GetSharedMediaQ
|
||||
return mediaType != "image" && mediaType != "video" && mediaType != "link";
|
||||
}
|
||||
|
||||
if (filterType == "media")
|
||||
{
|
||||
return mediaType == "image" || mediaType == "video";
|
||||
}
|
||||
|
||||
return true;
|
||||
}).ToList();
|
||||
|
||||
|
||||
@@ -16,21 +16,20 @@ public sealed record AddReactionCommand(
|
||||
|
||||
public sealed class AddReactionCommandHandler : ICommandHandler<AddReactionCommand>
|
||||
{
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IMessageReactionRepository _reactionRepository;
|
||||
private readonly IChatsUnitOfWork _unitOfWork;
|
||||
private readonly IHubContext<ChatHub> _hubContext;
|
||||
private readonly IUserDisplayNameProvider _displayNameProvider;
|
||||
private readonly ILogger<AddReactionCommandHandler> _logger;
|
||||
|
||||
public AddReactionCommandHandler(
|
||||
IMessageRepository messageRepository,
|
||||
|
||||
IMessageReactionRepository reactionRepository,
|
||||
IChatsUnitOfWork unitOfWork,
|
||||
IHubContext<ChatHub> hubContext,
|
||||
IUserDisplayNameProvider displayNameProvider,
|
||||
ILogger<AddReactionCommandHandler> logger)
|
||||
{
|
||||
_messageRepository = messageRepository;
|
||||
_reactionRepository = reactionRepository;
|
||||
_unitOfWork = unitOfWork;
|
||||
_hubContext = hubContext;
|
||||
_displayNameProvider = displayNameProvider;
|
||||
@@ -44,21 +43,8 @@ public sealed class AddReactionCommandHandler : ICommandHandler<AddReactionComma
|
||||
request.MessageId, request.UserId, request.Emoji, request.ChatId);
|
||||
|
||||
|
||||
var success = await _messageRepository.AddReactionAsync(
|
||||
request.MessageId,
|
||||
|
||||
request.UserId,
|
||||
|
||||
request.Emoji,
|
||||
|
||||
cancellationToken);
|
||||
|
||||
|
||||
if (!success)
|
||||
{
|
||||
_logger.LogWarning("AddReaction: Message not found {MessageId}", request.MessageId);
|
||||
return Result.Failure(ChatErrors.MessagesNotFound);
|
||||
}
|
||||
var reaction = new MessageReaction(request.MessageId, request.UserId, request.Emoji);
|
||||
await _reactionRepository.AddAsync(reaction, cancellationToken);
|
||||
|
||||
await _unitOfWork.SaveChangesAsync(cancellationToken);
|
||||
|
||||
|
||||
@@ -16,19 +16,18 @@ public sealed record RemoveReactionCommand(
|
||||
|
||||
public sealed class RemoveReactionCommandHandler : ICommandHandler<RemoveReactionCommand>
|
||||
{
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IMessageReactionRepository _reactionRepository;
|
||||
private readonly IChatsUnitOfWork _unitOfWork;
|
||||
private readonly IHubContext<ChatHub> _hubContext;
|
||||
private readonly ILogger<RemoveReactionCommandHandler> _logger;
|
||||
|
||||
public RemoveReactionCommandHandler(
|
||||
IMessageRepository messageRepository,
|
||||
|
||||
IMessageReactionRepository reactionRepository,
|
||||
IChatsUnitOfWork unitOfWork,
|
||||
IHubContext<ChatHub> hubContext,
|
||||
ILogger<RemoveReactionCommandHandler> logger)
|
||||
{
|
||||
_messageRepository = messageRepository;
|
||||
_reactionRepository = reactionRepository;
|
||||
_unitOfWork = unitOfWork;
|
||||
_hubContext = hubContext;
|
||||
_logger = logger;
|
||||
@@ -39,18 +38,12 @@ public sealed class RemoveReactionCommandHandler : ICommandHandler<RemoveReactio
|
||||
_logger.LogInformation("RemoveReaction: MessageId={MessageId}, UserId={UserId}, Emoji={Emoji}, ChatId={ChatId}",
|
||||
request.MessageId, request.UserId, request.Emoji, request.ChatId);
|
||||
|
||||
var success = await _messageRepository.RemoveReactionAsync(
|
||||
await _reactionRepository.RemoveAsync(
|
||||
request.MessageId,
|
||||
request.UserId,
|
||||
request.Emoji,
|
||||
cancellationToken);
|
||||
|
||||
if (!success)
|
||||
{
|
||||
_logger.LogWarning("RemoveReaction: Reaction not found");
|
||||
// Не возвращаем ошибку - реакция уже удалена или не существовала
|
||||
}
|
||||
|
||||
await _unitOfWork.SaveChangesAsync(cancellationToken);
|
||||
|
||||
_logger.LogInformation("RemoveReaction: Reaction removed from database");
|
||||
|
||||
@@ -5,27 +5,28 @@ using Knot.Modules.Chats.Application.Abstractions;
|
||||
|
||||
namespace Knot.Modules.Chats.Application.Messages.Read;
|
||||
|
||||
public sealed record ReadMessagesCommand(Guid ChatId, Guid UserId, List<Guid> MessageIds) : ICommand;
|
||||
public sealed record ReadMessagesCommand(Guid ChatId, Guid UserId, Guid LastReadMessageId, long LastReadSequenceId) : ICommand;
|
||||
|
||||
public sealed class ReadMessagesCommandHandler : ICommandHandler<ReadMessagesCommand>
|
||||
{
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IChatRepository _chatRepository;
|
||||
private readonly IChatsUnitOfWork _unitOfWork;
|
||||
|
||||
public ReadMessagesCommandHandler(IMessageRepository messageRepository, IChatsUnitOfWork unitOfWork)
|
||||
public ReadMessagesCommandHandler(IChatRepository chatRepository, IChatsUnitOfWork unitOfWork)
|
||||
{
|
||||
_messageRepository = messageRepository;
|
||||
_chatRepository = chatRepository;
|
||||
_unitOfWork = unitOfWork;
|
||||
}
|
||||
|
||||
public async Task<Result> Handle(ReadMessagesCommand request, CancellationToken cancellationToken)
|
||||
{
|
||||
if (request.MessageIds == null || !request.MessageIds.Any())
|
||||
{
|
||||
return Result.Success();
|
||||
}
|
||||
var chat = await _chatRepository.GetByIdAsync(request.ChatId, cancellationToken);
|
||||
if (chat == null) return Result.Failure(ChatErrors.NotFound);
|
||||
|
||||
await _messageRepository.AddReadReceiptsAsync(request.UserId, request.MessageIds, cancellationToken);
|
||||
var member = chat.Members.FirstOrDefault(m => m.UserId == request.UserId);
|
||||
if (member == null) return Result.Failure(ChatErrors.NotMember);
|
||||
|
||||
member.UpdateReadCursor(request.LastReadMessageId, request.LastReadSequenceId);
|
||||
|
||||
await _unitOfWork.SaveChangesAsync(cancellationToken);
|
||||
|
||||
|
||||
@@ -16,11 +16,13 @@ internal sealed class SearchMessagesQueryHandler : IQueryHandler<SearchMessagesQ
|
||||
{
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IUserDisplayNameProvider _userProvider;
|
||||
private readonly IMessageReactionRepository _reactionRepository;
|
||||
|
||||
public SearchMessagesQueryHandler(IMessageRepository messageRepository, IUserDisplayNameProvider userProvider)
|
||||
public SearchMessagesQueryHandler(IMessageRepository messageRepository, IUserDisplayNameProvider userProvider, IMessageReactionRepository reactionRepository)
|
||||
{
|
||||
_messageRepository = messageRepository;
|
||||
_userProvider = userProvider;
|
||||
_reactionRepository = reactionRepository;
|
||||
}
|
||||
|
||||
public async Task<Result<List<SearchMessageDto>>> Handle(SearchMessagesQuery request, CancellationToken cancellationToken)
|
||||
@@ -32,6 +34,10 @@ internal sealed class SearchMessagesQueryHandler : IQueryHandler<SearchMessagesQ
|
||||
userIds.AddRange(messages.Where(message => message.ForwardedFromId.HasValue).Select(message => message.ForwardedFromId!.Value));
|
||||
|
||||
var senders = await _userProvider.GetUsersInfoAsync(userIds.Distinct(), cancellationToken);
|
||||
|
||||
var messageIds = messages.Select(m => m.Id).ToList();
|
||||
var allReactions = await _reactionRepository.GetReactionsForMessagesAsync(messageIds, cancellationToken);
|
||||
var reactionsByMessage = allReactions.GroupBy(r => r.MessageId).ToDictionary(g => g.Key, g => g.ToList());
|
||||
|
||||
var result = messages.Select(message => new SearchMessageDto(
|
||||
message.Id,
|
||||
@@ -44,15 +50,16 @@ internal sealed class SearchMessagesQueryHandler : IQueryHandler<SearchMessagesQ
|
||||
message.IsEdited,
|
||||
message.IsDeleted,
|
||||
message.CreatedAt,
|
||||
message.SequenceId,
|
||||
message.ForwardedFromId,
|
||||
null,
|
||||
message.StoryId,
|
||||
message.StoryMediaUrl,
|
||||
message.StoryMediaType,
|
||||
message.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList(),
|
||||
message.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList(),
|
||||
senders.TryGetValue(message.SenderId, out var senderUser) ? new MessageSenderDto(senderUser.Id, senderUser.Username, senderUser.DisplayName, senderUser.Avatar) : new MessageSenderDto(message.SenderId, "unknown", "Unknown", null),
|
||||
message.Reactions.Select(reaction => new SimpleReactionDto(reaction.UserId, reaction.Emoji)).ToList(),
|
||||
message.ReadBy.Select(readReceipt => new ReadByDto(readReceipt.UserId)).ToList()
|
||||
reactionsByMessage.TryGetValue(message.Id, out var mr) ? mr.Select(reaction => new SimpleReactionDto(reaction.UserId, reaction.Emoji)).ToList() : new List<SimpleReactionDto>(),
|
||||
new List<ReadByDto>()
|
||||
)).ToList();
|
||||
|
||||
return Result.Success(result);
|
||||
|
||||
@@ -110,7 +110,15 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
|
||||
false);
|
||||
}
|
||||
|
||||
// 4. Сохраняем
|
||||
// 4. Последовательность сообщений High-Water Mark
|
||||
chat.IncrementSequenceId();
|
||||
message.SetSequenceId(chat.LastMessageSequenceId);
|
||||
|
||||
var senderMember = chat.Members.First(m => m.UserId == request.SenderId);
|
||||
senderMember.UpdateReadCursor(message.Id, message.SequenceId);
|
||||
senderMember.UpdateDeliveredCursor(message.Id);
|
||||
|
||||
// 5. Сохраняем
|
||||
_messageRepository.Add(message);
|
||||
await _unitOfWork.SaveChangesAsync(cancellationToken);
|
||||
|
||||
|
||||
@@ -35,6 +35,7 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImpor
|
||||
private readonly IMessageRepository _messageRepository;
|
||||
private readonly IFileStorageService _fileStorage;
|
||||
private readonly IHubContext<ChatHub> _hubContext;
|
||||
private readonly IMessageReactionRepository _reactionRepository;
|
||||
|
||||
public ExecuteImportCommandHandler(
|
||||
ISender sender,
|
||||
@@ -42,7 +43,8 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImpor
|
||||
IChatRepository chatRepository,
|
||||
IMessageRepository messageRepository,
|
||||
IFileStorageService fileStorage,
|
||||
IHubContext<ChatHub> hubContext)
|
||||
IHubContext<ChatHub> hubContext,
|
||||
IMessageReactionRepository reactionRepository)
|
||||
{
|
||||
_sender = sender;
|
||||
_uow = uow;
|
||||
@@ -50,6 +52,7 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImpor
|
||||
_messageRepository = messageRepository;
|
||||
_fileStorage = fileStorage;
|
||||
_hubContext = hubContext;
|
||||
_reactionRepository = reactionRepository;
|
||||
}
|
||||
|
||||
public async Task<Result<ExecuteImportResponseDto>> Handle(ExecuteImportCommand request, CancellationToken cancellationToken)
|
||||
@@ -349,13 +352,25 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImpor
|
||||
}
|
||||
|
||||
bool isJoined = fromNameNode == null;
|
||||
bool isMediaOnly = string.IsNullOrEmpty(content) && forwardedNode == null && replyToId == null;
|
||||
bool hasMedia = mediaNodes != null && mediaNodes.Count > 0;
|
||||
Message? targetMessage = null;
|
||||
|
||||
if (isJoined && isMediaOnly && lastSavedMessage is MediaMessage && Math.Abs((createdAt - lastSavedMessage.CreatedAt).TotalSeconds) <= 60 && lastSavedMessage.SenderId == senderGuid)
|
||||
// Check if we should combine this message with the previous one
|
||||
// We combine if: it's joined AND it's just media/text within 60s AND same sender
|
||||
// Even if it's forwarded, Telegram exports media groups as joined forwarded messages.
|
||||
bool shouldCombine = isJoined && lastSavedMessage is MediaMessage
|
||||
&& Math.Abs((createdAt - lastSavedMessage.CreatedAt).TotalSeconds) <= 60
|
||||
&& lastSavedMessage.SenderId == senderGuid
|
||||
&& lastSavedMessage.ForwardedFromId == forwardedFromId;
|
||||
|
||||
if (shouldCombine)
|
||||
{
|
||||
targetMessage = lastSavedMessage;
|
||||
if (targetMessage is MediaMessage mm && !string.IsNullOrEmpty(content) && content != mm.Content)
|
||||
{
|
||||
// If the joined message has text (caption), append it
|
||||
mm.AppendImportedCaption(content);
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -378,33 +393,51 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImpor
|
||||
|
||||
if (hasMedia)
|
||||
{
|
||||
var seenMedia = new HashSet<string>();
|
||||
var validMediaExtracted = new List<(string href, string finalMType, string cType)>();
|
||||
|
||||
foreach (var mediaNode in mediaNodes!)
|
||||
{
|
||||
string? href = mediaNode.GetAttribute("href") ?? mediaNode.GetAttribute("src");
|
||||
if (!string.IsNullOrEmpty(href) && !href.StartsWith("http"))
|
||||
{
|
||||
var zipPath = baseDir + href.Replace("\\", "/");
|
||||
var zipEntry = archive.GetEntry(zipPath);
|
||||
if (zipEntry != null)
|
||||
if (!seenMedia.Add(href)) continue;
|
||||
|
||||
var types = GetMediaTypes(href);
|
||||
var finalMType = types.mType;
|
||||
if (mediaNode.ClassName?.Contains("animated") == true)
|
||||
{
|
||||
using var ms = new MemoryStream();
|
||||
using var zipfs = zipEntry.Open();
|
||||
await zipfs.CopyToAsync(ms, cancellationToken);
|
||||
ms.Position = 0;
|
||||
finalMType = "image";
|
||||
}
|
||||
|
||||
validMediaExtracted.Add((href, finalMType, types.cType));
|
||||
}
|
||||
}
|
||||
|
||||
var types = GetMediaTypes(href);
|
||||
var finalMType = types.mType;
|
||||
if (mediaNode.ClassName?.Contains("animated") == true)
|
||||
{
|
||||
finalMType = "image";
|
||||
}
|
||||
// В Telegram экспорте если в одном .message блоке есть и видео, и картинка - картинка это просто миниатюра (thumbnail).
|
||||
// Реальные альбомы идут отдельными .message div'ами c классом joined.
|
||||
// Поэтому мы просто удаляем картинку, чтобы она не дублировалась как отдельный файл в галерее!
|
||||
if (validMediaExtracted.Any(m => m.finalMType == "video") && validMediaExtracted.Any(m => m.finalMType == "image"))
|
||||
{
|
||||
validMediaExtracted.RemoveAll(m => m.finalMType == "image");
|
||||
}
|
||||
|
||||
var parsedType = Enum.TryParse<MediaType>(finalMType, true, out var mTypeEnum) ? mTypeEnum : MediaType.File;
|
||||
if (targetMessage is MediaMessage mm)
|
||||
{
|
||||
var fileId = await _fileStorage.UploadFileAsync(ms, Path.GetFileName(href), types.cType);
|
||||
mm.AddMedia(parsedType, $"/api/files/{fileId}", Path.GetFileName(href), zipEntry.Length);
|
||||
}
|
||||
foreach (var mediaTuple in validMediaExtracted)
|
||||
{
|
||||
var zipPath = baseDir + mediaTuple.href.Replace("\\", "/");
|
||||
var zipEntry = archive.GetEntry(zipPath);
|
||||
if (zipEntry != null)
|
||||
{
|
||||
using var ms = new MemoryStream();
|
||||
using var zipfs = zipEntry.Open();
|
||||
await zipfs.CopyToAsync(ms, cancellationToken);
|
||||
ms.Position = 0;
|
||||
|
||||
var parsedType = Enum.TryParse<MediaType>(mediaTuple.finalMType, true, out var mTypeEnum) ? mTypeEnum : MediaType.File;
|
||||
if (targetMessage is MediaMessage mm)
|
||||
{
|
||||
var fileId = await _fileStorage.UploadFileAsync(ms, Path.GetFileName(mediaTuple.href), mediaTuple.cType);
|
||||
mm.AddMedia(parsedType, $"/api/files/{fileId}", Path.GetFileName(mediaTuple.href), zipEntry.Length);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -425,18 +458,23 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImpor
|
||||
foreach (var userpicNode in userpicNodes)
|
||||
{
|
||||
var title = userpicNode.GetAttribute("title")?.Trim();
|
||||
if (!string.IsNullOrEmpty(title) && request.Mapping.TryGetValue(title, out var rUserId) && rUserId != Guid.Empty)
|
||||
if (!string.IsNullOrEmpty(title) && request.Mapping.TryGetValue(title, out var rUserId) && rUserId != Guid.Empty && targetMessage != null)
|
||||
{
|
||||
targetMessage.AddReaction(rUserId, emoji);
|
||||
var reaction = new MessageReaction(targetMessage.Id, rUserId, emoji);
|
||||
await _reactionRepository.AddAsync(reaction, cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (targetMessage != lastSavedMessage && (!string.IsNullOrEmpty(content) || targetMessage.Media.Any()))
|
||||
if (targetMessage != null && targetMessage != lastSavedMessage && (!string.IsNullOrEmpty(content) || targetMessage.Media.Any()))
|
||||
{
|
||||
_messageRepository.Add(targetMessage);
|
||||
lastSavedMessage = targetMessage;
|
||||
importedCount++;
|
||||
}
|
||||
else if (shouldCombine && lastSavedMessage != null)
|
||||
{
|
||||
await _messageRepository.UpdateAsync(lastSavedMessage, cancellationToken);
|
||||
}
|
||||
}
|
||||
catch { /* ignore single message parse error */ }
|
||||
}
|
||||
|
||||
@@ -38,6 +38,7 @@ public static class DependencyInjection
|
||||
services.AddScoped<IChatsUnitOfWork>(sp => sp.GetRequiredService<ChatsDbContext>());
|
||||
services.AddScoped<IChatRepository, ChatRepository>();
|
||||
services.AddScoped<IMessageRepository, MessageRepository>();
|
||||
services.AddScoped<IMessageReactionRepository, MessageReactionRepository>();
|
||||
|
||||
// MediatR
|
||||
services.AddMediatR(config =>
|
||||
|
||||
@@ -35,6 +35,7 @@ public sealed class Chat : AggregateRoot<Guid>
|
||||
public string? Description { get; private set; }
|
||||
public string? Avatar { get; private set; }
|
||||
public DateTime CreatedAt { get; private set; }
|
||||
public long LastMessageSequenceId { get; private set; }
|
||||
|
||||
private readonly List<ChatMember> _members = new();
|
||||
public IReadOnlyCollection<ChatMember> Members => _members.AsReadOnly();
|
||||
@@ -103,6 +104,11 @@ public sealed class Chat : AggregateRoot<Guid>
|
||||
public void UpdateDescription(string? description) => Description = description;
|
||||
|
||||
public void UpdateAvatar(string? avatarUrl) => Avatar = avatarUrl;
|
||||
|
||||
public long IncrementSequenceId()
|
||||
{
|
||||
return ++LastMessageSequenceId;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -116,6 +122,10 @@ public sealed class ChatMember : Entity<Guid>
|
||||
public DateTime JoinedAt { get; private set; }
|
||||
public bool IsPinned { get; private set; }
|
||||
public bool IsMuted { get; private set; }
|
||||
|
||||
public Guid? LastReadMessageId { get; private set; }
|
||||
public long LastReadSequenceId { get; private set; }
|
||||
public Guid? LastDeliveredMessageId { get; private set; }
|
||||
|
||||
// For EF Core
|
||||
private ChatMember() : base(Guid.Empty) { Role = "member"; }
|
||||
@@ -129,4 +139,18 @@ public sealed class ChatMember : Entity<Guid>
|
||||
}
|
||||
|
||||
public void TogglePin() => IsPinned = !IsPinned;
|
||||
|
||||
public void UpdateReadCursor(Guid messageId, long sequenceId)
|
||||
{
|
||||
if (sequenceId > LastReadSequenceId)
|
||||
{
|
||||
LastReadMessageId = messageId;
|
||||
LastReadSequenceId = sequenceId;
|
||||
}
|
||||
}
|
||||
|
||||
public void UpdateDeliveredCursor(Guid messageId)
|
||||
{
|
||||
LastDeliveredMessageId = messageId;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,6 +9,8 @@ public static class ChatErrors
|
||||
public static readonly Error ImportExpired = new Error("Import.Expired", "Session not found or expired");
|
||||
public static readonly Error ImportMissing = new Error("Import.Missing", "ZIP file lost");
|
||||
public static readonly Error ChatNotFound = new Error("Chat.NotFound", "Chat not found or access denied");
|
||||
public static readonly Error NotFound = new Error("Chat.NotFound", "Chat not found"); // Alias
|
||||
public static readonly Error NotMember = new Error("Chat.NotMember", "You are not a member of this chat");
|
||||
public static readonly Error ChatsForbidden = new Error("Chats.Forbidden", "Вы не являетесь участником этого чата.");
|
||||
public static readonly Error MessagesNotFound = new Error("Messages.NotFound", "Message not found.");
|
||||
public static readonly Error ChatsNotFound = new Error("Chats.NotFound", "Чат не найден.");
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Knot.Modules.Chats.Domain;
|
||||
|
||||
public interface IMessageReactionRepository
|
||||
{
|
||||
Task AddAsync(MessageReaction reaction, CancellationToken cancellationToken);
|
||||
Task RemoveAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken);
|
||||
Task<List<MessageReaction>> GetReactionsForMessageAsync(Guid messageId, CancellationToken cancellationToken);
|
||||
Task<List<MessageReaction>> GetReactionsForMessagesAsync(IEnumerable<Guid> messageIds, CancellationToken cancellationToken);
|
||||
}
|
||||
@@ -9,11 +9,9 @@ public interface IMessageRepository
|
||||
Task<List<Message>> GetChatMessagesAsync(Guid chatId, int limit, int offset, CancellationToken cancellationToken);
|
||||
Task<Message?> GetLatestChatMessageAsync(Guid chatId, CancellationToken cancellationToken);
|
||||
Task<List<Message>> SearchMessagesAsync(string query, Guid? chatId, Guid requestingUserId, CancellationToken cancellationToken);
|
||||
Task AddReadReceiptsAsync(Guid userId, List<Guid> messageIds, CancellationToken cancellationToken);
|
||||
Task<bool> AddReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken);
|
||||
Task<bool> RemoveReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken);
|
||||
|
||||
Task<List<Message>> GetChatMessagesCursorAsync(Guid chatId, DateTime? cursor, int limit, CancellationToken cancellationToken);
|
||||
Task<Message?> GetLastStoryMessageAsync(Guid chatId, Guid storyId, CancellationToken cancellationToken);
|
||||
Task<int> GetUnreadCountAsync(Guid chatId, Guid userId, CancellationToken cancellationToken);
|
||||
|
||||
Task UpdateAsync(Message message, CancellationToken cancellationToken);
|
||||
}
|
||||
|
||||
@@ -44,6 +44,18 @@ public class MediaMessage : Message
|
||||
_media.Add(new Domain.Media(Id, type.ToString().ToLower(), url, filename, size));
|
||||
}
|
||||
|
||||
public void AppendImportedCaption(string additionalCaption)
|
||||
{
|
||||
if (string.IsNullOrEmpty(Caption))
|
||||
{
|
||||
Caption = additionalCaption;
|
||||
}
|
||||
else
|
||||
{
|
||||
Caption += "\n" + additionalCaption;
|
||||
}
|
||||
}
|
||||
|
||||
public void Edit(string newCaption)
|
||||
{
|
||||
Caption = newCaption;
|
||||
|
||||
@@ -13,7 +13,12 @@ public abstract class Message : AggregateRoot<Guid>
|
||||
public Guid ChatId { get; protected set; }
|
||||
public Guid SenderId { get; protected set; }
|
||||
public DateTime CreatedAt { get; protected set; }
|
||||
public long SequenceId { get; protected set; }
|
||||
|
||||
public void SetSequenceId(long sequenceId)
|
||||
{
|
||||
SequenceId = sequenceId;
|
||||
}
|
||||
// ================== Опциональные метаданные (общего назначения) ==================
|
||||
public Guid? ReplyToId { get; protected set; }
|
||||
public Guid? ForwardedFromId { get; protected set; }
|
||||
@@ -36,15 +41,9 @@ public abstract class Message : AggregateRoot<Guid>
|
||||
public bool IsDeletedForUser(Guid userId) => _deletedFor.Exists(d => d.UserId == userId);
|
||||
|
||||
// ================== Связанные коллекции (общего назначения) ==================
|
||||
protected List<ReadReceipt> _readBy = new();
|
||||
public IReadOnlyCollection<ReadReceipt> ReadBy => _readBy.AsReadOnly();
|
||||
|
||||
protected List<DeletedMessage> _deletedFor = new();
|
||||
public IReadOnlyCollection<DeletedMessage> DeletedFor => _deletedFor.AsReadOnly();
|
||||
|
||||
protected List<Reaction> _reactions = new();
|
||||
public IReadOnlyCollection<Reaction> Reactions => _reactions.AsReadOnly();
|
||||
|
||||
// ================== Инфраструктурный конструктор EF ==================
|
||||
protected Message() : base(Guid.Empty) { }
|
||||
|
||||
@@ -72,28 +71,9 @@ public abstract class Message : AggregateRoot<Guid>
|
||||
public bool HasState(MessageState state) => (State & state) == state;
|
||||
|
||||
// ================== Общие операции ==================
|
||||
public void AddReaction(Guid userId, string emoji)
|
||||
{
|
||||
var existing = _reactions.Find(r => r.UserId == userId && r.Emoji == emoji);
|
||||
if (existing == null)
|
||||
{
|
||||
_reactions.Add(new Reaction(Id, userId, emoji));
|
||||
}
|
||||
}
|
||||
|
||||
public void RemoveReaction(Guid userId, string emoji)
|
||||
{
|
||||
var existing = _reactions.Find(r => r.UserId == userId && r.Emoji == emoji);
|
||||
if (existing != null)
|
||||
{
|
||||
_reactions.Remove(existing);
|
||||
}
|
||||
}
|
||||
|
||||
public virtual void Delete()
|
||||
{
|
||||
AddState(MessageState.IsDeleted);
|
||||
_reactions.Clear();
|
||||
}
|
||||
|
||||
public void DeleteForUser(Guid userId)
|
||||
|
||||
29
backend/src/Modules/Chats/Domain/MessageReaction.cs
Normal file
29
backend/src/Modules/Chats/Domain/MessageReaction.cs
Normal file
@@ -0,0 +1,29 @@
|
||||
using System;
|
||||
using Knot.Shared.Kernel;
|
||||
|
||||
namespace Knot.Modules.Chats.Domain;
|
||||
|
||||
/// <summary>
|
||||
/// Агрегат/Сущность реакции для сообщения. Вынесен в отдельную коллекцию
|
||||
/// для бесконечного масштабирования и чистоты DDD (Approach 3).
|
||||
/// </summary>
|
||||
public sealed class MessageReaction : AggregateRoot<Guid>
|
||||
{
|
||||
public Guid MessageId { get; private set; }
|
||||
public Guid UserId { get; private set; }
|
||||
public string Emoji { get; private set; }
|
||||
public DateTime CreatedAt { get; private set; }
|
||||
|
||||
private MessageReaction() : base(Guid.Empty)
|
||||
{
|
||||
Emoji = default!;
|
||||
}
|
||||
|
||||
public MessageReaction(Guid messageId, Guid userId, string emoji) : base(Guid.NewGuid())
|
||||
{
|
||||
MessageId = messageId;
|
||||
UserId = userId;
|
||||
Emoji = emoji;
|
||||
CreatedAt = DateTime.UtcNow;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
using MongoDB.Driver;
|
||||
using Knot.Modules.Chats.Domain;
|
||||
|
||||
namespace Knot.Modules.Chats.Infrastructure.Persistence;
|
||||
|
||||
public sealed class MessageReactionRepository : IMessageReactionRepository
|
||||
{
|
||||
private readonly IMongoCollection<MessageReaction> _reactions;
|
||||
|
||||
public MessageReactionRepository(IMongoDatabase mongoDatabase)
|
||||
{
|
||||
_reactions = mongoDatabase.GetCollection<MessageReaction>("message_reactions");
|
||||
|
||||
// Ensure index for fast querying by message
|
||||
var indexKeysDefinition = Builders<MessageReaction>.IndexKeys.Ascending(r => r.MessageId);
|
||||
_reactions.Indexes.CreateOne(new CreateIndexModel<MessageReaction>(indexKeysDefinition));
|
||||
}
|
||||
|
||||
public async Task AddAsync(MessageReaction reaction, CancellationToken cancellationToken)
|
||||
{
|
||||
// Уникальный индекс или фильтр, чтобы не дублировать
|
||||
var filter = Builders<MessageReaction>.Filter.And(
|
||||
Builders<MessageReaction>.Filter.Eq(r => r.MessageId, reaction.MessageId),
|
||||
Builders<MessageReaction>.Filter.Eq(r => r.UserId, reaction.UserId),
|
||||
Builders<MessageReaction>.Filter.Eq(r => r.Emoji, reaction.Emoji)
|
||||
);
|
||||
|
||||
// Используем ReplaceOptions.IsUpsert = true для идемпотентности (нет гонок)
|
||||
await _reactions.ReplaceOneAsync(filter, reaction, new ReplaceOptions { IsUpsert = true }, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task RemoveAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken)
|
||||
{
|
||||
var filter = Builders<MessageReaction>.Filter.And(
|
||||
Builders<MessageReaction>.Filter.Eq(r => r.MessageId, messageId),
|
||||
Builders<MessageReaction>.Filter.Eq(r => r.UserId, userId),
|
||||
Builders<MessageReaction>.Filter.Eq(r => r.Emoji, emoji)
|
||||
);
|
||||
|
||||
await _reactions.DeleteOneAsync(filter, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<List<MessageReaction>> GetReactionsForMessageAsync(Guid messageId, CancellationToken cancellationToken)
|
||||
{
|
||||
return await _reactions.Find(r => r.MessageId == messageId).ToListAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<List<MessageReaction>> GetReactionsForMessagesAsync(IEnumerable<Guid> messageIds, CancellationToken cancellationToken)
|
||||
{
|
||||
var filter = Builders<MessageReaction>.Filter.In(r => r.MessageId, messageIds);
|
||||
return await _reactions.Find(filter).ToListAsync(cancellationToken);
|
||||
}
|
||||
}
|
||||
@@ -102,81 +102,7 @@ public sealed class MessageRepository : IMessageRepository
|
||||
.ToListAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task AddReadReceiptsAsync(Guid userId, List<Guid> messageIds, CancellationToken cancellationToken)
|
||||
{
|
||||
var filter = Builders<Message>.Filter.In(m => m.Id, messageIds);
|
||||
var messages = await _messages.Find(filter).ToListAsync(cancellationToken);
|
||||
if (!messages.Any()) return;
|
||||
|
||||
var chatId = messages.First().ChatId;
|
||||
var maxDate = messages.Max(m => m.CreatedAt);
|
||||
|
||||
var notReadFilter = Builders<Message>.Filter.Not(
|
||||
Builders<Message>.Filter.ElemMatch<ReadReceipt>(
|
||||
"ReadBy",
|
||||
Builders<ReadReceipt>.Filter.Eq(r => r.UserId, userId)
|
||||
)
|
||||
);
|
||||
|
||||
var finalFilter = Builders<Message>.Filter.And(
|
||||
Builders<Message>.Filter.Eq(m => m.ChatId, chatId),
|
||||
Builders<Message>.Filter.Ne(m => m.SenderId, userId),
|
||||
Builders<Message>.Filter.Lte(m => m.CreatedAt, maxDate),
|
||||
notReadFilter
|
||||
);
|
||||
|
||||
var unreadMessagesToMark = await _messages.Find(finalFilter).ToListAsync(cancellationToken);
|
||||
|
||||
var writes = new List<WriteModel<Message>>();
|
||||
foreach (var msg in unreadMessagesToMark)
|
||||
{
|
||||
var receipt = new ReadReceipt(msg.Id, userId);
|
||||
var pushUpdate = Builders<Message>.Update.Push("ReadBy", receipt);
|
||||
var updateModel = new UpdateOneModel<Message>(Builders<Message>.Filter.Eq(m => m.Id, msg.Id), pushUpdate);
|
||||
writes.Add(updateModel);
|
||||
}
|
||||
|
||||
if (writes.Any())
|
||||
{
|
||||
await _messages.BulkWriteAsync(writes, cancellationToken: cancellationToken);
|
||||
}
|
||||
}
|
||||
|
||||
public async Task<bool> AddReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken)
|
||||
{
|
||||
var filter = Builders<Message>.Filter.Eq(m => m.Id, messageId);
|
||||
var msg = await _messages.Find(filter).FirstOrDefaultAsync(cancellationToken);
|
||||
if (msg == null) return false;
|
||||
|
||||
if (msg.Reactions.Any(r => r.UserId == userId && r.Emoji == emoji))
|
||||
{
|
||||
return true;
|
||||
}
|
||||
|
||||
var reaction = new Reaction(messageId, userId, emoji);
|
||||
var update = Builders<Message>.Update.Push("Reactions", reaction);
|
||||
await _messages.UpdateOneAsync(filter, update, cancellationToken: cancellationToken);
|
||||
return true;
|
||||
}
|
||||
|
||||
public async Task<bool> RemoveReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken)
|
||||
{
|
||||
var filter = Builders<Message>.Filter.Eq(m => m.Id, messageId);
|
||||
var msg = await _messages.Find(filter).FirstOrDefaultAsync(cancellationToken);
|
||||
if (msg == null) return false;
|
||||
|
||||
var reaction = msg.Reactions.FirstOrDefault(r => r.UserId == userId && r.Emoji == emoji);
|
||||
if (reaction == null) return false;
|
||||
|
||||
var update = Builders<Message>.Update.PullFilter("Reactions",
|
||||
Builders<BsonDocument>.Filter.And(
|
||||
Builders<BsonDocument>.Filter.Eq("UserId", userId),
|
||||
Builders<BsonDocument>.Filter.Eq("Emoji", emoji)
|
||||
));
|
||||
|
||||
await _messages.UpdateOneAsync(filter, update, cancellationToken: cancellationToken);
|
||||
return true;
|
||||
}
|
||||
|
||||
public async Task<Message?> GetLastStoryMessageAsync(Guid chatId, Guid storyId, CancellationToken cancellationToken)
|
||||
{
|
||||
@@ -191,23 +117,7 @@ public sealed class MessageRepository : IMessageRepository
|
||||
.FirstOrDefaultAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<int> GetUnreadCountAsync(Guid chatId, Guid userId, CancellationToken cancellationToken)
|
||||
{
|
||||
var notReadFilter = Builders<Message>.Filter.Not(
|
||||
Builders<Message>.Filter.ElemMatch<ReadReceipt>(
|
||||
"ReadBy",
|
||||
Builders<ReadReceipt>.Filter.Eq(r => r.UserId, userId)
|
||||
)
|
||||
);
|
||||
|
||||
var finalFilter = Builders<Message>.Filter.And(
|
||||
Builders<Message>.Filter.Eq(m => m.ChatId, chatId),
|
||||
Builders<Message>.Filter.Ne(m => m.SenderId, userId),
|
||||
notReadFilter
|
||||
);
|
||||
|
||||
return (int)await _messages.CountDocumentsAsync(finalFilter, cancellationToken: cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UpdateAsync(Message message, CancellationToken cancellationToken)
|
||||
{
|
||||
|
||||
@@ -29,8 +29,6 @@ public static class MongoDbMapConfigurator
|
||||
{
|
||||
cm.AutoMap();
|
||||
cm.MapField("_deletedFor").SetElementName("DeletedFor");
|
||||
cm.MapField("_readBy").SetElementName("ReadBy");
|
||||
cm.MapField("_reactions").SetElementName("Reactions");
|
||||
cm.MapProperty(c => c.Content).SetSerializer(new EncryptedStringSerializer());
|
||||
cm.MapProperty(c => c.Quote).SetSerializer(new EncryptedStringSerializer());
|
||||
cm.SetIsRootClass(true);
|
||||
@@ -58,8 +56,8 @@ public static class MongoDbMapConfigurator
|
||||
});
|
||||
|
||||
BsonClassMap.RegisterClassMap<DeletedMessage>(cm => cm.AutoMap());
|
||||
BsonClassMap.RegisterClassMap<ReadReceipt>(cm => cm.AutoMap());
|
||||
BsonClassMap.RegisterClassMap<Reaction>(cm => cm.AutoMap());
|
||||
BsonClassMap.RegisterClassMap<MessageReaction>(cm => cm.AutoMap());
|
||||
|
||||
|
||||
BsonClassMap.RegisterClassMap<Media>(cm =>
|
||||
{
|
||||
|
||||
@@ -118,26 +118,19 @@ public sealed class ChatHub : Hub
|
||||
[HubMethodName("read_messages")]
|
||||
public async Task ReadMessages(ReadMessagesRequest request)
|
||||
{
|
||||
if (request.MessageIds != null && request.MessageIds.Any())
|
||||
if (request.LastReadMessageId != Guid.Empty && request.LastReadSequenceId > 0)
|
||||
{
|
||||
var parsedIds = request.MessageIds
|
||||
.Select(id => Guid.TryParse(id, out var parsed) ? parsed : Guid.Empty)
|
||||
.Where(id => id != Guid.Empty)
|
||||
.ToList();
|
||||
|
||||
if (parsedIds.Any())
|
||||
{
|
||||
var command = new ReadMessagesCommand(
|
||||
request.ChatId, _userContext.UserId, parsedIds);
|
||||
await _sender.Send(command);
|
||||
}
|
||||
var command = new ReadMessagesCommand(
|
||||
request.ChatId, _userContext.UserId, request.LastReadMessageId, request.LastReadSequenceId);
|
||||
await _sender.Send(command);
|
||||
}
|
||||
|
||||
await Clients.Group(request.ChatId.ToString()).SendAsync("messages_read", new
|
||||
{
|
||||
ChatId = request.ChatId.ToString(),
|
||||
UserId = _userContext.UserId,
|
||||
MessageIds = request.MessageIds ?? new List<string>()
|
||||
LastReadMessageId = request.LastReadMessageId,
|
||||
LastReadSequenceId = request.LastReadSequenceId
|
||||
});
|
||||
}
|
||||
|
||||
@@ -644,7 +637,7 @@ public sealed class ChatHub : Hub
|
||||
Guid? ReplyToId = null,
|
||||
string? Quote = null,
|
||||
Guid? ForwardedFromId = null);
|
||||
public record ReadMessagesRequest(Guid ChatId, List<string>? MessageIds);
|
||||
public record ReadMessagesRequest(Guid ChatId, Guid LastReadMessageId, long LastReadSequenceId);
|
||||
public record CallOfferRequest(string TargetUserId, object Offer, string CallType, string? ChatId);
|
||||
public record CallAnswerRequest(string TargetUserId, object Answer);
|
||||
public record TargetUserRequest(string TargetUserId);
|
||||
|
||||
111
backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.Designer.cs
generated
Normal file
111
backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.Designer.cs
generated
Normal file
@@ -0,0 +1,111 @@
|
||||
// <auto-generated />
|
||||
using System;
|
||||
using Knot.Modules.Chats.Infrastructure.Persistence;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.EntityFrameworkCore.Infrastructure;
|
||||
using Microsoft.EntityFrameworkCore.Migrations;
|
||||
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
|
||||
using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
|
||||
|
||||
#nullable disable
|
||||
|
||||
namespace Knot.Modules.Chats.Migrations
|
||||
{
|
||||
[DbContext(typeof(ChatsDbContext))]
|
||||
[Migration("20260320185319_AddHighWaterMark")]
|
||||
partial class AddHighWaterMark
|
||||
{
|
||||
/// <inheritdoc />
|
||||
protected override void BuildTargetModel(ModelBuilder modelBuilder)
|
||||
{
|
||||
#pragma warning disable 612, 618
|
||||
modelBuilder
|
||||
.HasDefaultSchema("chats")
|
||||
.HasAnnotation("ProductVersion", "10.0.4")
|
||||
.HasAnnotation("Relational:MaxIdentifierLength", 63);
|
||||
|
||||
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
|
||||
|
||||
modelBuilder.Entity("Knot.Modules.Chats.Domain.Chat", b =>
|
||||
{
|
||||
b.Property<Guid>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b.Property<string>("Avatar")
|
||||
.HasColumnType("text");
|
||||
|
||||
b.Property<DateTime>("CreatedAt")
|
||||
.HasColumnType("timestamp with time zone");
|
||||
|
||||
b.Property<string>("Description")
|
||||
.HasColumnType("text");
|
||||
|
||||
b.Property<long>("LastMessageSequenceId")
|
||||
.HasColumnType("bigint");
|
||||
|
||||
b.Property<string>("Name")
|
||||
.HasColumnType("text");
|
||||
|
||||
b.Property<string>("Type")
|
||||
.IsRequired()
|
||||
.HasColumnType("text");
|
||||
|
||||
b.HasKey("Id");
|
||||
|
||||
b.ToTable("Chats", "chats");
|
||||
});
|
||||
|
||||
modelBuilder.Entity("Knot.Modules.Chats.Domain.Chat", b =>
|
||||
{
|
||||
b.OwnsMany("Knot.Modules.Chats.Domain.ChatMember", "Members", b1 =>
|
||||
{
|
||||
b1.Property<Guid>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b1.Property<Guid>("ChatId")
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b1.Property<bool>("IsMuted")
|
||||
.HasColumnType("boolean");
|
||||
|
||||
b1.Property<bool>("IsPinned")
|
||||
.HasColumnType("boolean");
|
||||
|
||||
b1.Property<DateTime>("JoinedAt")
|
||||
.HasColumnType("timestamp with time zone");
|
||||
|
||||
b1.Property<Guid?>("LastDeliveredMessageId")
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b1.Property<Guid?>("LastReadMessageId")
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b1.Property<long>("LastReadSequenceId")
|
||||
.HasColumnType("bigint");
|
||||
|
||||
b1.Property<string>("Role")
|
||||
.IsRequired()
|
||||
.HasColumnType("text");
|
||||
|
||||
b1.Property<Guid>("UserId")
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b1.HasKey("Id");
|
||||
|
||||
b1.HasIndex("ChatId", "UserId")
|
||||
.IsUnique();
|
||||
|
||||
b1.ToTable("ChatMembers", "chats");
|
||||
|
||||
b1.WithOwner()
|
||||
.HasForeignKey("ChatId");
|
||||
});
|
||||
|
||||
b.Navigation("Members");
|
||||
});
|
||||
#pragma warning restore 612, 618
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,69 @@
|
||||
using System;
|
||||
using Microsoft.EntityFrameworkCore.Migrations;
|
||||
|
||||
#nullable disable
|
||||
|
||||
namespace Knot.Modules.Chats.Migrations
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public partial class AddHighWaterMark : Migration
|
||||
{
|
||||
/// <inheritdoc />
|
||||
protected override void Up(MigrationBuilder migrationBuilder)
|
||||
{
|
||||
migrationBuilder.AddColumn<long>(
|
||||
name: "LastMessageSequenceId",
|
||||
schema: "chats",
|
||||
table: "Chats",
|
||||
type: "bigint",
|
||||
nullable: false,
|
||||
defaultValue: 0L);
|
||||
|
||||
migrationBuilder.AddColumn<Guid>(
|
||||
name: "LastDeliveredMessageId",
|
||||
schema: "chats",
|
||||
table: "ChatMembers",
|
||||
type: "uuid",
|
||||
nullable: true);
|
||||
|
||||
migrationBuilder.AddColumn<Guid>(
|
||||
name: "LastReadMessageId",
|
||||
schema: "chats",
|
||||
table: "ChatMembers",
|
||||
type: "uuid",
|
||||
nullable: true);
|
||||
|
||||
migrationBuilder.AddColumn<long>(
|
||||
name: "LastReadSequenceId",
|
||||
schema: "chats",
|
||||
table: "ChatMembers",
|
||||
type: "bigint",
|
||||
nullable: false,
|
||||
defaultValue: 0L);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
protected override void Down(MigrationBuilder migrationBuilder)
|
||||
{
|
||||
migrationBuilder.DropColumn(
|
||||
name: "LastMessageSequenceId",
|
||||
schema: "chats",
|
||||
table: "Chats");
|
||||
|
||||
migrationBuilder.DropColumn(
|
||||
name: "LastDeliveredMessageId",
|
||||
schema: "chats",
|
||||
table: "ChatMembers");
|
||||
|
||||
migrationBuilder.DropColumn(
|
||||
name: "LastReadMessageId",
|
||||
schema: "chats",
|
||||
table: "ChatMembers");
|
||||
|
||||
migrationBuilder.DropColumn(
|
||||
name: "LastReadSequenceId",
|
||||
schema: "chats",
|
||||
table: "ChatMembers");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -38,6 +38,9 @@ namespace Knot.Modules.Chats.Migrations
|
||||
b.Property<string>("Description")
|
||||
.HasColumnType("text");
|
||||
|
||||
b.Property<long>("LastMessageSequenceId")
|
||||
.HasColumnType("bigint");
|
||||
|
||||
b.Property<string>("Name")
|
||||
.HasColumnType("text");
|
||||
|
||||
@@ -70,6 +73,15 @@ namespace Knot.Modules.Chats.Migrations
|
||||
b1.Property<DateTime>("JoinedAt")
|
||||
.HasColumnType("timestamp with time zone");
|
||||
|
||||
b1.Property<Guid?>("LastDeliveredMessageId")
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b1.Property<Guid?>("LastReadMessageId")
|
||||
.HasColumnType("uuid");
|
||||
|
||||
b1.Property<long>("LastReadSequenceId")
|
||||
.HasColumnType("bigint");
|
||||
|
||||
b1.Property<string>("Role")
|
||||
.IsRequired()
|
||||
.HasColumnType("text");
|
||||
|
||||
Reference in New Issue
Block a user