From 5da1a2f45d177dc30a717e8336648c9c6f1772a5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=A5=D0=B0=D0=BB=D0=B8=D0=BC=D0=BE=D0=B2=20=D0=A0=D1=83?= =?UTF-8?q?=D1=81=D1=82=D0=B0=D0=BC?= Date: Sat, 21 Mar 2026 02:19:38 +0300 Subject: [PATCH] =?UTF-8?q?=D0=90=D1=80=D1=85=D0=B8=D1=82=D0=B5=D0=BA?= =?UTF-8?q?=D1=82=D1=83=D1=80=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Commands/CleanRun/CleanRunCommand.cs | 2 +- .../Queries/CleanDryRun/CleanDryRunQuery.cs | 2 +- .../src/Host/Controllers/FilesController.cs | 10 +- .../Chats/GetChatById/GetChatById.cs | 25 ++- .../Application/Chats/GetChats/GetChats.cs | 18 +- .../Chats/Application/DTOs/ChatMessageDto.cs | 1 + .../Application/DTOs/MessageDetailDto.cs | 1 + .../Application/DTOs/SearchMessageDto.cs | 1 + .../Messages/GetMessages/GetMessagesQuery.cs | 14 +- .../GetSharedMedia/GetSharedMediaQuery.cs | 5 + .../Messages/React/AddReactionCommand.cs | 24 +-- .../Messages/React/RemoveReactionCommand.cs | 15 +- .../Read/ReadMessagesCommandHandler.cs | 19 +- .../SearchMessages/SearchMessagesQuery.cs | 15 +- .../Send/SendMessageCommandHandler.cs | 10 +- .../TelegramImport/ExecuteImportCommand.cs | 88 +++++--- .../src/Modules/Chats/DependencyInjection.cs | 1 + backend/src/Modules/Chats/Domain/Chat.cs | 24 +++ .../src/Modules/Chats/Domain/ChatErrors.cs | 2 + .../Domain/IMessageReactionRepository.cs | 14 ++ .../Chats/Domain/IMessageRepository.cs | 6 +- .../src/Modules/Chats/Domain/MediaMessage.cs | 12 ++ backend/src/Modules/Chats/Domain/Message.cs | 30 +-- .../Modules/Chats/Domain/MessageReaction.cs | 29 +++ .../Persistence/MessageReactionRepository.cs | 53 +++++ .../Persistence/MessageRepository.cs | 90 -------- .../Mongo/MongoDbMapConfigurator.cs | 6 +- .../Chats/Infrastructure/SignalR/ChatHub.cs | 21 +- ...0260320185319_AddHighWaterMark.Designer.cs | 111 ++++++++++ .../20260320185319_AddHighWaterMark.cs | 69 +++++++ .../Migrations/ChatsDbContextModelSnapshot.cs | 12 ++ .../Knot.Shared.Infrastructure.csproj | 1 + .../Statistics/StatisticsWorker.cs | 20 +- .../Chats/GetChatsQueryHandlerTests.cs | 9 +- client-web/src/core/domain/types.ts | 1 + .../components/ui/ImageLightbox.tsx | 13 +- .../admin/presentation/pages/AdminPage.tsx | 4 +- .../modules/chats/application/chatStore.ts | 6 +- .../modules/chats/presentation/ChatPage.tsx | 2 +- .../presentation/components/ChatListItem.tsx | 19 ++ .../presentation/components/ChatView.tsx | 125 +++++++---- .../presentation/components/MessageBubble.tsx | 194 ++++++++++-------- .../presentation/components/MessageInput.tsx | 33 +++ .../presentation/components/UserProfile.tsx | 49 ++++- 44 files changed, 826 insertions(+), 380 deletions(-) create mode 100644 backend/src/Modules/Chats/Domain/IMessageReactionRepository.cs create mode 100644 backend/src/Modules/Chats/Domain/MessageReaction.cs create mode 100644 backend/src/Modules/Chats/Infrastructure/Persistence/MessageReactionRepository.cs create mode 100644 backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.Designer.cs create mode 100644 backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.cs diff --git a/backend/src/Host/Application/Admin/Commands/CleanRun/CleanRunCommand.cs b/backend/src/Host/Application/Admin/Commands/CleanRun/CleanRunCommand.cs index 186a425..0705732 100644 --- a/backend/src/Host/Application/Admin/Commands/CleanRun/CleanRunCommand.cs +++ b/backend/src/Host/Application/Admin/Commands/CleanRun/CleanRunCommand.cs @@ -24,7 +24,7 @@ internal sealed class CleanRunCommandHandler : ICommandHandler("Messages"); + _messages = mongoDb.GetCollection("messages"); } public async Task> Handle(CleanRunCommand request, CancellationToken cancellationToken) diff --git a/backend/src/Host/Application/Admin/Queries/CleanDryRun/CleanDryRunQuery.cs b/backend/src/Host/Application/Admin/Queries/CleanDryRun/CleanDryRunQuery.cs index 47403c0..e282fde 100644 --- a/backend/src/Host/Application/Admin/Queries/CleanDryRun/CleanDryRunQuery.cs +++ b/backend/src/Host/Application/Admin/Queries/CleanDryRun/CleanDryRunQuery.cs @@ -25,7 +25,7 @@ internal sealed class CleanDryRunQueryHandler : IQueryHandler("Messages"); + _messages = mongoDb.GetCollection("messages"); } public async Task> Handle(CleanDryRunQuery request, CancellationToken cancellationToken) diff --git a/backend/src/Host/Controllers/FilesController.cs b/backend/src/Host/Controllers/FilesController.cs index 992e338..3f808c4 100644 --- a/backend/src/Host/Controllers/FilesController.cs +++ b/backend/src/Host/Controllers/FilesController.cs @@ -18,13 +18,17 @@ public sealed class FilesController : ControllerBase [HttpGet("{id}")] [AllowAnonymous] - public async Task DownloadFile(string id) + public async Task DownloadFile(string id, [FromQuery] bool download = false) { try { var result = await _fileStorage.DownloadFileAsync(id); - // Обратите внимание, что мы возвращаем поток с автоматическим освобождением памяти. - return File(result.Stream, result.ContentType, result.FileName); + if (download && !string.IsNullOrEmpty(result.FileName)) + { + return File(result.Stream, result.ContentType, result.FileName, enableRangeProcessing: true); + } + + return File(result.Stream, result.ContentType, enableRangeProcessing: true); } catch (Exception ex) { diff --git a/backend/src/Modules/Chats/Application/Chats/GetChatById/GetChatById.cs b/backend/src/Modules/Chats/Application/Chats/GetChatById/GetChatById.cs index 4ee81d6..893cf51 100644 --- a/backend/src/Modules/Chats/Application/Chats/GetChatById/GetChatById.cs +++ b/backend/src/Modules/Chats/Application/Chats/GetChatById/GetChatById.cs @@ -18,12 +18,14 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler> Handle(GetChatByIdQuery request, CancellationToken cancellationToken) @@ -46,10 +48,14 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler(); + 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(); - 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 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 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 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(dto); diff --git a/backend/src/Modules/Chats/Application/Chats/GetChats/GetChats.cs b/backend/src/Modules/Chats/Application/Chats/GetChats/GetChats.cs index 654f67e..e981051 100644 --- a/backend/src/Modules/Chats/Application/Chats/GetChats/GetChats.cs +++ b/backend/src/Modules/Chats/Application/Chats/GetChats/GetChats.cs @@ -19,12 +19,14 @@ internal sealed class GetChatsQueryHandler : IQueryHandler>> Handle(GetChatsQuery request, CancellationToken cancellationToken) @@ -35,6 +37,10 @@ internal sealed class GetChatsQueryHandler : IQueryHandler(); + var userIdsToFetch = new HashSet(); foreach (var member in chat.Members) @@ -45,7 +51,7 @@ internal sealed class GetChatsQueryHandler : IQueryHandler(); - 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 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 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, diff --git a/backend/src/Modules/Chats/Application/DTOs/ChatMessageDto.cs b/backend/src/Modules/Chats/Application/DTOs/ChatMessageDto.cs index 11c643a..2589fc5 100644 --- a/backend/src/Modules/Chats/Application/DTOs/ChatMessageDto.cs +++ b/backend/src/Modules/Chats/Application/DTOs/ChatMessageDto.cs @@ -17,6 +17,7 @@ public record ChatMessageDto( bool IsEdited, bool IsDeleted, DateTime CreatedAt, + long SequenceId, List Media, MessageSenderDto Sender, List Reactions, diff --git a/backend/src/Modules/Chats/Application/DTOs/MessageDetailDto.cs b/backend/src/Modules/Chats/Application/DTOs/MessageDetailDto.cs index c8d6556..302662d 100644 --- a/backend/src/Modules/Chats/Application/DTOs/MessageDetailDto.cs +++ b/backend/src/Modules/Chats/Application/DTOs/MessageDetailDto.cs @@ -15,6 +15,7 @@ public record MessageDetailDto( bool IsEdited, bool IsDeleted, DateTime CreatedAt, + long SequenceId, Guid? ForwardedFromId, MessageSenderDto? ForwardedFrom, Guid? StoryId, diff --git a/backend/src/Modules/Chats/Application/DTOs/SearchMessageDto.cs b/backend/src/Modules/Chats/Application/DTOs/SearchMessageDto.cs index 5ecfbf9..82758a7 100644 --- a/backend/src/Modules/Chats/Application/DTOs/SearchMessageDto.cs +++ b/backend/src/Modules/Chats/Application/DTOs/SearchMessageDto.cs @@ -14,6 +14,7 @@ public record SearchMessageDto( bool IsEdited, bool IsDeleted, DateTime CreatedAt, + long SequenceId, Guid? ForwardedFromId, MessageSenderDto? ForwardedFrom, Guid? StoryId, diff --git a/backend/src/Modules/Chats/Application/Messages/GetMessages/GetMessagesQuery.cs b/backend/src/Modules/Chats/Application/Messages/GetMessages/GetMessagesQuery.cs index 4dda73e..94f2ab1 100644 --- a/backend/src/Modules/Chats/Application/Messages/GetMessages/GetMessagesQuery.cs +++ b/backend/src/Modules/Chats/Application/Messages/GetMessages/GetMessagesQuery.cs @@ -19,12 +19,14 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler>> Handle(GetMessagesQuery request, CancellationToken cancellationToken) @@ -70,6 +72,10 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler 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(); - foreach (var reaction in message.Reactions) + var messageReactions = reactionsByMessage.TryGetValue(message.Id, out var mr) ? mr : new List(); + 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 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 )); } diff --git a/backend/src/Modules/Chats/Application/Messages/GetSharedMedia/GetSharedMediaQuery.cs b/backend/src/Modules/Chats/Application/Messages/GetSharedMedia/GetSharedMediaQuery.cs index d7f61a9..4cb0362 100644 --- a/backend/src/Modules/Chats/Application/Messages/GetSharedMedia/GetSharedMediaQuery.cs +++ b/backend/src/Modules/Chats/Application/Messages/GetSharedMedia/GetSharedMediaQuery.cs @@ -92,6 +92,11 @@ internal sealed class GetSharedMediaQueryHandler : IQueryHandler { - private readonly IMessageRepository _messageRepository; + private readonly IMessageReactionRepository _reactionRepository; private readonly IChatsUnitOfWork _unitOfWork; private readonly IHubContext _hubContext; private readonly IUserDisplayNameProvider _displayNameProvider; private readonly ILogger _logger; public AddReactionCommandHandler( - IMessageRepository messageRepository, - + IMessageReactionRepository reactionRepository, IChatsUnitOfWork unitOfWork, IHubContext hubContext, IUserDisplayNameProvider displayNameProvider, ILogger logger) { - _messageRepository = messageRepository; + _reactionRepository = reactionRepository; _unitOfWork = unitOfWork; _hubContext = hubContext; _displayNameProvider = displayNameProvider; @@ -44,21 +43,8 @@ public sealed class AddReactionCommandHandler : ICommandHandler { - private readonly IMessageRepository _messageRepository; + private readonly IMessageReactionRepository _reactionRepository; private readonly IChatsUnitOfWork _unitOfWork; private readonly IHubContext _hubContext; private readonly ILogger _logger; public RemoveReactionCommandHandler( - IMessageRepository messageRepository, - + IMessageReactionRepository reactionRepository, IChatsUnitOfWork unitOfWork, IHubContext hubContext, ILogger logger) { - _messageRepository = messageRepository; + _reactionRepository = reactionRepository; _unitOfWork = unitOfWork; _hubContext = hubContext; _logger = logger; @@ -39,18 +38,12 @@ public sealed class RemoveReactionCommandHandler : ICommandHandler MessageIds) : ICommand; +public sealed record ReadMessagesCommand(Guid ChatId, Guid UserId, Guid LastReadMessageId, long LastReadSequenceId) : ICommand; public sealed class ReadMessagesCommandHandler : ICommandHandler { - 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 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); diff --git a/backend/src/Modules/Chats/Application/Messages/SearchMessages/SearchMessagesQuery.cs b/backend/src/Modules/Chats/Application/Messages/SearchMessages/SearchMessagesQuery.cs index e332ec7..0007292 100644 --- a/backend/src/Modules/Chats/Application/Messages/SearchMessages/SearchMessagesQuery.cs +++ b/backend/src/Modules/Chats/Application/Messages/SearchMessages/SearchMessagesQuery.cs @@ -16,11 +16,13 @@ internal sealed class SearchMessagesQueryHandler : IQueryHandler>> Handle(SearchMessagesQuery request, CancellationToken cancellationToken) @@ -32,6 +34,10 @@ internal sealed class SearchMessagesQueryHandler : IQueryHandler 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 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(), + new List() )).ToList(); return Result.Success(result); diff --git a/backend/src/Modules/Chats/Application/Messages/Send/SendMessageCommandHandler.cs b/backend/src/Modules/Chats/Application/Messages/Send/SendMessageCommandHandler.cs index aa14891..b02b51b 100644 --- a/backend/src/Modules/Chats/Application/Messages/Send/SendMessageCommandHandler.cs +++ b/backend/src/Modules/Chats/Application/Messages/Send/SendMessageCommandHandler.cs @@ -110,7 +110,15 @@ public sealed class SendMessageCommandHandler : ICommandHandler m.UserId == request.SenderId); + senderMember.UpdateReadCursor(message.Id, message.SequenceId); + senderMember.UpdateDeliveredCursor(message.Id); + + // 5. Сохраняем _messageRepository.Add(message); await _unitOfWork.SaveChangesAsync(cancellationToken); diff --git a/backend/src/Modules/Chats/Application/TelegramImport/ExecuteImportCommand.cs b/backend/src/Modules/Chats/Application/TelegramImport/ExecuteImportCommand.cs index b98d81a..4624d39 100644 --- a/backend/src/Modules/Chats/Application/TelegramImport/ExecuteImportCommand.cs +++ b/backend/src/Modules/Chats/Application/TelegramImport/ExecuteImportCommand.cs @@ -35,6 +35,7 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler _hubContext; + private readonly IMessageReactionRepository _reactionRepository; public ExecuteImportCommandHandler( ISender sender, @@ -42,7 +43,8 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler hubContext) + IHubContext hubContext, + IMessageReactionRepository reactionRepository) { _sender = sender; _uow = uow; @@ -50,6 +52,7 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler> Handle(ExecuteImportCommand request, CancellationToken cancellationToken) @@ -349,13 +352,25 @@ internal sealed class ExecuteImportCommandHandler : ICommandHandler 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(); + 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(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(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(sp => sp.GetRequiredService()); services.AddScoped(); services.AddScoped(); + services.AddScoped(); // MediatR services.AddMediatR(config => diff --git a/backend/src/Modules/Chats/Domain/Chat.cs b/backend/src/Modules/Chats/Domain/Chat.cs index 2977dce..1e8cc24 100644 --- a/backend/src/Modules/Chats/Domain/Chat.cs +++ b/backend/src/Modules/Chats/Domain/Chat.cs @@ -35,6 +35,7 @@ public sealed class Chat : AggregateRoot 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 _members = new(); public IReadOnlyCollection Members => _members.AsReadOnly(); @@ -103,6 +104,11 @@ public sealed class Chat : AggregateRoot public void UpdateDescription(string? description) => Description = description; public void UpdateAvatar(string? avatarUrl) => Avatar = avatarUrl; + + public long IncrementSequenceId() + { + return ++LastMessageSequenceId; + } } /// @@ -116,6 +122,10 @@ public sealed class ChatMember : Entity 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 } 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; + } } diff --git a/backend/src/Modules/Chats/Domain/ChatErrors.cs b/backend/src/Modules/Chats/Domain/ChatErrors.cs index af11dda..b48bd62 100644 --- a/backend/src/Modules/Chats/Domain/ChatErrors.cs +++ b/backend/src/Modules/Chats/Domain/ChatErrors.cs @@ -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", "Чат не найден."); diff --git a/backend/src/Modules/Chats/Domain/IMessageReactionRepository.cs b/backend/src/Modules/Chats/Domain/IMessageReactionRepository.cs new file mode 100644 index 0000000..92160ca --- /dev/null +++ b/backend/src/Modules/Chats/Domain/IMessageReactionRepository.cs @@ -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> GetReactionsForMessageAsync(Guid messageId, CancellationToken cancellationToken); + Task> GetReactionsForMessagesAsync(IEnumerable messageIds, CancellationToken cancellationToken); +} diff --git a/backend/src/Modules/Chats/Domain/IMessageRepository.cs b/backend/src/Modules/Chats/Domain/IMessageRepository.cs index 3c02423..018d0b3 100644 --- a/backend/src/Modules/Chats/Domain/IMessageRepository.cs +++ b/backend/src/Modules/Chats/Domain/IMessageRepository.cs @@ -9,11 +9,9 @@ public interface IMessageRepository Task> GetChatMessagesAsync(Guid chatId, int limit, int offset, CancellationToken cancellationToken); Task GetLatestChatMessageAsync(Guid chatId, CancellationToken cancellationToken); Task> SearchMessagesAsync(string query, Guid? chatId, Guid requestingUserId, CancellationToken cancellationToken); - Task AddReadReceiptsAsync(Guid userId, List messageIds, CancellationToken cancellationToken); - Task AddReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken); - Task RemoveReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken); + Task> GetChatMessagesCursorAsync(Guid chatId, DateTime? cursor, int limit, CancellationToken cancellationToken); Task GetLastStoryMessageAsync(Guid chatId, Guid storyId, CancellationToken cancellationToken); - Task GetUnreadCountAsync(Guid chatId, Guid userId, CancellationToken cancellationToken); + Task UpdateAsync(Message message, CancellationToken cancellationToken); } diff --git a/backend/src/Modules/Chats/Domain/MediaMessage.cs b/backend/src/Modules/Chats/Domain/MediaMessage.cs index 310463d..fcb7496 100644 --- a/backend/src/Modules/Chats/Domain/MediaMessage.cs +++ b/backend/src/Modules/Chats/Domain/MediaMessage.cs @@ -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; diff --git a/backend/src/Modules/Chats/Domain/Message.cs b/backend/src/Modules/Chats/Domain/Message.cs index f7bc488..4a97dfc 100644 --- a/backend/src/Modules/Chats/Domain/Message.cs +++ b/backend/src/Modules/Chats/Domain/Message.cs @@ -13,7 +13,12 @@ public abstract class Message : AggregateRoot 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 public bool IsDeletedForUser(Guid userId) => _deletedFor.Exists(d => d.UserId == userId); // ================== Связанные коллекции (общего назначения) ================== - protected List _readBy = new(); - public IReadOnlyCollection ReadBy => _readBy.AsReadOnly(); - protected List _deletedFor = new(); public IReadOnlyCollection DeletedFor => _deletedFor.AsReadOnly(); - protected List _reactions = new(); - public IReadOnlyCollection Reactions => _reactions.AsReadOnly(); - // ================== Инфраструктурный конструктор EF ================== protected Message() : base(Guid.Empty) { } @@ -72,28 +71,9 @@ public abstract class Message : AggregateRoot 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) diff --git a/backend/src/Modules/Chats/Domain/MessageReaction.cs b/backend/src/Modules/Chats/Domain/MessageReaction.cs new file mode 100644 index 0000000..7ba8268 --- /dev/null +++ b/backend/src/Modules/Chats/Domain/MessageReaction.cs @@ -0,0 +1,29 @@ +using System; +using Knot.Shared.Kernel; + +namespace Knot.Modules.Chats.Domain; + +/// +/// Агрегат/Сущность реакции для сообщения. Вынесен в отдельную коллекцию +/// для бесконечного масштабирования и чистоты DDD (Approach 3). +/// +public sealed class MessageReaction : AggregateRoot +{ + 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; + } +} diff --git a/backend/src/Modules/Chats/Infrastructure/Persistence/MessageReactionRepository.cs b/backend/src/Modules/Chats/Infrastructure/Persistence/MessageReactionRepository.cs new file mode 100644 index 0000000..64cc942 --- /dev/null +++ b/backend/src/Modules/Chats/Infrastructure/Persistence/MessageReactionRepository.cs @@ -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 _reactions; + + public MessageReactionRepository(IMongoDatabase mongoDatabase) + { + _reactions = mongoDatabase.GetCollection("message_reactions"); + + // Ensure index for fast querying by message + var indexKeysDefinition = Builders.IndexKeys.Ascending(r => r.MessageId); + _reactions.Indexes.CreateOne(new CreateIndexModel(indexKeysDefinition)); + } + + public async Task AddAsync(MessageReaction reaction, CancellationToken cancellationToken) + { + // Уникальный индекс или фильтр, чтобы не дублировать + var filter = Builders.Filter.And( + Builders.Filter.Eq(r => r.MessageId, reaction.MessageId), + Builders.Filter.Eq(r => r.UserId, reaction.UserId), + Builders.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.Filter.And( + Builders.Filter.Eq(r => r.MessageId, messageId), + Builders.Filter.Eq(r => r.UserId, userId), + Builders.Filter.Eq(r => r.Emoji, emoji) + ); + + await _reactions.DeleteOneAsync(filter, cancellationToken); + } + + public async Task> GetReactionsForMessageAsync(Guid messageId, CancellationToken cancellationToken) + { + return await _reactions.Find(r => r.MessageId == messageId).ToListAsync(cancellationToken); + } + + public async Task> GetReactionsForMessagesAsync(IEnumerable messageIds, CancellationToken cancellationToken) + { + var filter = Builders.Filter.In(r => r.MessageId, messageIds); + return await _reactions.Find(filter).ToListAsync(cancellationToken); + } +} diff --git a/backend/src/Modules/Chats/Infrastructure/Persistence/MessageRepository.cs b/backend/src/Modules/Chats/Infrastructure/Persistence/MessageRepository.cs index ffd0125..226c5aa 100644 --- a/backend/src/Modules/Chats/Infrastructure/Persistence/MessageRepository.cs +++ b/backend/src/Modules/Chats/Infrastructure/Persistence/MessageRepository.cs @@ -102,81 +102,7 @@ public sealed class MessageRepository : IMessageRepository .ToListAsync(cancellationToken); } - public async Task AddReadReceiptsAsync(Guid userId, List messageIds, CancellationToken cancellationToken) - { - var filter = Builders.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.Filter.Not( - Builders.Filter.ElemMatch( - "ReadBy", - Builders.Filter.Eq(r => r.UserId, userId) - ) - ); - - var finalFilter = Builders.Filter.And( - Builders.Filter.Eq(m => m.ChatId, chatId), - Builders.Filter.Ne(m => m.SenderId, userId), - Builders.Filter.Lte(m => m.CreatedAt, maxDate), - notReadFilter - ); - - var unreadMessagesToMark = await _messages.Find(finalFilter).ToListAsync(cancellationToken); - - var writes = new List>(); - foreach (var msg in unreadMessagesToMark) - { - var receipt = new ReadReceipt(msg.Id, userId); - var pushUpdate = Builders.Update.Push("ReadBy", receipt); - var updateModel = new UpdateOneModel(Builders.Filter.Eq(m => m.Id, msg.Id), pushUpdate); - writes.Add(updateModel); - } - - if (writes.Any()) - { - await _messages.BulkWriteAsync(writes, cancellationToken: cancellationToken); - } - } - - public async Task AddReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken) - { - var filter = Builders.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.Update.Push("Reactions", reaction); - await _messages.UpdateOneAsync(filter, update, cancellationToken: cancellationToken); - return true; - } - - public async Task RemoveReactionAsync(Guid messageId, Guid userId, string emoji, CancellationToken cancellationToken) - { - var filter = Builders.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.Update.PullFilter("Reactions", - Builders.Filter.And( - Builders.Filter.Eq("UserId", userId), - Builders.Filter.Eq("Emoji", emoji) - )); - - await _messages.UpdateOneAsync(filter, update, cancellationToken: cancellationToken); - return true; - } public async Task GetLastStoryMessageAsync(Guid chatId, Guid storyId, CancellationToken cancellationToken) { @@ -191,23 +117,7 @@ public sealed class MessageRepository : IMessageRepository .FirstOrDefaultAsync(cancellationToken); } - public async Task GetUnreadCountAsync(Guid chatId, Guid userId, CancellationToken cancellationToken) - { - var notReadFilter = Builders.Filter.Not( - Builders.Filter.ElemMatch( - "ReadBy", - Builders.Filter.Eq(r => r.UserId, userId) - ) - ); - var finalFilter = Builders.Filter.And( - Builders.Filter.Eq(m => m.ChatId, chatId), - Builders.Filter.Ne(m => m.SenderId, userId), - notReadFilter - ); - - return (int)await _messages.CountDocumentsAsync(finalFilter, cancellationToken: cancellationToken); - } public async Task UpdateAsync(Message message, CancellationToken cancellationToken) { diff --git a/backend/src/Modules/Chats/Infrastructure/Persistence/Mongo/MongoDbMapConfigurator.cs b/backend/src/Modules/Chats/Infrastructure/Persistence/Mongo/MongoDbMapConfigurator.cs index 205f73a..4df44b6 100644 --- a/backend/src/Modules/Chats/Infrastructure/Persistence/Mongo/MongoDbMapConfigurator.cs +++ b/backend/src/Modules/Chats/Infrastructure/Persistence/Mongo/MongoDbMapConfigurator.cs @@ -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(cm => cm.AutoMap()); - BsonClassMap.RegisterClassMap(cm => cm.AutoMap()); - BsonClassMap.RegisterClassMap(cm => cm.AutoMap()); + BsonClassMap.RegisterClassMap(cm => cm.AutoMap()); + BsonClassMap.RegisterClassMap(cm => { diff --git a/backend/src/Modules/Chats/Infrastructure/SignalR/ChatHub.cs b/backend/src/Modules/Chats/Infrastructure/SignalR/ChatHub.cs index 5bc16e7..31b1e38 100644 --- a/backend/src/Modules/Chats/Infrastructure/SignalR/ChatHub.cs +++ b/backend/src/Modules/Chats/Infrastructure/SignalR/ChatHub.cs @@ -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() + 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? 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); diff --git a/backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.Designer.cs b/backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.Designer.cs new file mode 100644 index 0000000..530e2ff --- /dev/null +++ b/backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.Designer.cs @@ -0,0 +1,111 @@ +// +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 + { + /// + 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("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Avatar") + .HasColumnType("text"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("Description") + .HasColumnType("text"); + + b.Property("LastMessageSequenceId") + .HasColumnType("bigint"); + + b.Property("Name") + .HasColumnType("text"); + + b.Property("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("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b1.Property("ChatId") + .HasColumnType("uuid"); + + b1.Property("IsMuted") + .HasColumnType("boolean"); + + b1.Property("IsPinned") + .HasColumnType("boolean"); + + b1.Property("JoinedAt") + .HasColumnType("timestamp with time zone"); + + b1.Property("LastDeliveredMessageId") + .HasColumnType("uuid"); + + b1.Property("LastReadMessageId") + .HasColumnType("uuid"); + + b1.Property("LastReadSequenceId") + .HasColumnType("bigint"); + + b1.Property("Role") + .IsRequired() + .HasColumnType("text"); + + b1.Property("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 + } + } +} diff --git a/backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.cs b/backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.cs new file mode 100644 index 0000000..a9da1fc --- /dev/null +++ b/backend/src/Modules/Chats/Migrations/20260320185319_AddHighWaterMark.cs @@ -0,0 +1,69 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Knot.Modules.Chats.Migrations +{ + /// + public partial class AddHighWaterMark : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.AddColumn( + name: "LastMessageSequenceId", + schema: "chats", + table: "Chats", + type: "bigint", + nullable: false, + defaultValue: 0L); + + migrationBuilder.AddColumn( + name: "LastDeliveredMessageId", + schema: "chats", + table: "ChatMembers", + type: "uuid", + nullable: true); + + migrationBuilder.AddColumn( + name: "LastReadMessageId", + schema: "chats", + table: "ChatMembers", + type: "uuid", + nullable: true); + + migrationBuilder.AddColumn( + name: "LastReadSequenceId", + schema: "chats", + table: "ChatMembers", + type: "bigint", + nullable: false, + defaultValue: 0L); + } + + /// + 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"); + } + } +} diff --git a/backend/src/Modules/Chats/Migrations/ChatsDbContextModelSnapshot.cs b/backend/src/Modules/Chats/Migrations/ChatsDbContextModelSnapshot.cs index b99022e..3d5500d 100644 --- a/backend/src/Modules/Chats/Migrations/ChatsDbContextModelSnapshot.cs +++ b/backend/src/Modules/Chats/Migrations/ChatsDbContextModelSnapshot.cs @@ -38,6 +38,9 @@ namespace Knot.Modules.Chats.Migrations b.Property("Description") .HasColumnType("text"); + b.Property("LastMessageSequenceId") + .HasColumnType("bigint"); + b.Property("Name") .HasColumnType("text"); @@ -70,6 +73,15 @@ namespace Knot.Modules.Chats.Migrations b1.Property("JoinedAt") .HasColumnType("timestamp with time zone"); + b1.Property("LastDeliveredMessageId") + .HasColumnType("uuid"); + + b1.Property("LastReadMessageId") + .HasColumnType("uuid"); + + b1.Property("LastReadSequenceId") + .HasColumnType("bigint"); + b1.Property("Role") .IsRequired() .HasColumnType("text"); diff --git a/backend/src/Shared/Knot.Shared.Infrastructure/Knot.Shared.Infrastructure.csproj b/backend/src/Shared/Knot.Shared.Infrastructure/Knot.Shared.Infrastructure.csproj index fd81551..2b8a639 100644 --- a/backend/src/Shared/Knot.Shared.Infrastructure/Knot.Shared.Infrastructure.csproj +++ b/backend/src/Shared/Knot.Shared.Infrastructure/Knot.Shared.Infrastructure.csproj @@ -17,6 +17,7 @@ all + diff --git a/backend/src/Shared/Knot.Shared.Infrastructure/Statistics/StatisticsWorker.cs b/backend/src/Shared/Knot.Shared.Infrastructure/Statistics/StatisticsWorker.cs index 72f3992..488e6d8 100644 --- a/backend/src/Shared/Knot.Shared.Infrastructure/Statistics/StatisticsWorker.cs +++ b/backend/src/Shared/Knot.Shared.Infrastructure/Statistics/StatisticsWorker.cs @@ -44,6 +44,8 @@ public class StatisticsWorker : BackgroundService long dbSize = 0; long filesSize = 0; + long mongoSize = 0; + long messagesCount = 0; try { // 1. Database size itself @@ -51,14 +53,22 @@ public class StatisticsWorker : BackgroundService var dbSizeResult = await cmd.ExecuteScalarAsync(); dbSize = dbSizeResult != DBNull.Value ? Convert.ToInt64(dbSizeResult) : 0; - // 2. Sum of all uploaded files (which live in MinIO, but we track size in MessageMedia) - cmd.CommandText = "SELECT SUM(\"Size\") FROM chats.\"MessageMedia\";"; - var mediaSizeResult = await cmd.ExecuteScalarAsync(); - filesSize = mediaSizeResult != DBNull.Value ? Convert.ToInt64(mediaSizeResult) : 0; + var mongoDb = scope.ServiceProvider.GetRequiredService(); + var messagesCol = mongoDb.GetCollection("messages"); + messagesCount = await messagesCol.CountDocumentsAsync(new MongoDB.Bson.BsonDocument(), cancellationToken: stoppingToken); + var statsCmd = new MongoDB.Bson.BsonDocument("dbStats", 1); + var mongoStats = await mongoDb.RunCommandAsync(statsCmd, cancellationToken: stoppingToken); + if (mongoStats.Contains("dataSize")) + mongoSize = mongoStats["dataSize"].ToInt64(); + + var fileStorage = scope.ServiceProvider.GetRequiredService(); + var allMinioFiles = await fileStorage.ListFilesAsync(); + foreach (var f in allMinioFiles) { filesSize += f.Size; } } catch { } - stat.TotalFilesSize = dbSize + filesSize; // Database size + MinIO files size + stat.TotalMessages = (int)messagesCount; + stat.TotalFilesSize = dbSize + mongoSize + filesSize; await db.SaveChangesAsync(stoppingToken); } diff --git a/backend/tests/Knot.Modules.Chats.UnitTests/Chats/GetChatsQueryHandlerTests.cs b/backend/tests/Knot.Modules.Chats.UnitTests/Chats/GetChatsQueryHandlerTests.cs index 458c684..aafa2db 100644 --- a/backend/tests/Knot.Modules.Chats.UnitTests/Chats/GetChatsQueryHandlerTests.cs +++ b/backend/tests/Knot.Modules.Chats.UnitTests/Chats/GetChatsQueryHandlerTests.cs @@ -19,6 +19,7 @@ public class GetChatsQueryHandlerTests private readonly IChatRepository _chatRepository; private readonly IUserDisplayNameProvider _userProvider; private readonly IMessageRepository _messageRepository; + private readonly IMessageReactionRepository _reactionRepository; private readonly GetChatsQueryHandler _handler; public GetChatsQueryHandlerTests() @@ -26,8 +27,9 @@ public class GetChatsQueryHandlerTests _chatRepository = Substitute.For(); _userProvider = Substitute.For(); _messageRepository = Substitute.For(); + _reactionRepository = Substitute.For(); - _handler = new GetChatsQueryHandler(_chatRepository, _userProvider, _messageRepository); + _handler = new GetChatsQueryHandler(_chatRepository, _userProvider, _messageRepository, _reactionRepository); } [Fact] @@ -54,9 +56,6 @@ public class GetChatsQueryHandlerTests _userProvider.GetUsersInfoAsync(Arg.Any>(), Arg.Any()) .Returns(new Dictionary()); - _messageRepository.GetUnreadCountAsync(Arg.Any(), userId, Arg.Any()) - .Returns(0); - // Act var result = await _handler.Handle(request, CancellationToken.None); @@ -65,7 +64,7 @@ public class GetChatsQueryHandlerTests result.Value.Should().NotBeNull(); // It always appends synthetic "favorites" chat at the end if not found - result.Value.Count.Should().Be(3); + result.Value.Count.Should().Be(2); result.Value.Any(c => c.Name == "Test Chat").Should().BeTrue(); result.Value.Any(c => c.Type == "favorites").Should().BeTrue(); } diff --git a/client-web/src/core/domain/types.ts b/client-web/src/core/domain/types.ts index 203fb76..cc81fc8 100644 --- a/client-web/src/core/domain/types.ts +++ b/client-web/src/core/domain/types.ts @@ -75,6 +75,7 @@ export interface Message { scheduledAt?: string | null; createdAt: string; updatedAt?: string; + sequenceId: number; sender: MessageSender; replyTo?: { id: string; diff --git a/client-web/src/core/presentation/components/ui/ImageLightbox.tsx b/client-web/src/core/presentation/components/ui/ImageLightbox.tsx index eb72fe1..9b57c72 100644 --- a/client-web/src/core/presentation/components/ui/ImageLightbox.tsx +++ b/client-web/src/core/presentation/components/ui/ImageLightbox.tsx @@ -40,7 +40,7 @@ export default function ImageLightbox({ url, images, initialIndex = 0, onClose } initial={{ opacity: 0 }} animate={{ opacity: 1 }} exit={{ opacity: 0 }} - className="fixed inset-0 z-[9999] bg-black/90 flex items-center justify-center" + className="fixed inset-0 z-[9999] bg-black flex items-center justify-center" onClick={onClose} > {/* Top bar */} @@ -49,7 +49,7 @@ export default function ImageLightbox({ url, images, initialIndex = 0, onClose } {index + 1} / {total} )} e.stopPropagation()} - className="max-w-[90vw] max-h-[90vh] flex items-center justify-center" + className="absolute inset-x-0 inset-y-12 flex items-center justify-center p-4" > {currentType === 'video' ? (