Исправлен механизм прочтения сообщений

This commit is contained in:
Халимов Рустам
2026-04-20 21:44:39 +03:00
parent 0f593e52e0
commit 83ed328dd5
12 changed files with 224 additions and 162 deletions
@@ -19,13 +19,16 @@ public abstract class Message : AggregateRoot<Guid>
public bool IsEdited => HasState(MessageState.IsEdited);
public bool IsDeleted => HasState(MessageState.IsDeleted);
protected List<DeletedMessage> _deletedFor = new();
public IReadOnlyCollection<DeletedMessage> DeletedFor => _deletedFor.AsReadOnly();
protected List<Guid> _readByUsers = new();
public IReadOnlyCollection<Guid> ReadByUsers => _readByUsers.AsReadOnly();
protected Message() : base(Guid.Empty) { }
protected Message(Guid id, Guid chatId, Guid senderId, Guid? replyToId, Guid? forwardedFromId, DateTime createdAt, bool isImported)
protected Message(Guid id, Guid chatId, Guid senderId, Guid? replyToId, Guid? forwardedFromId, DateTime createdAt, bool isImported)
: base(id)
{
ChatId = chatId;
@@ -42,15 +45,25 @@ public abstract class Message : AggregateRoot<Guid>
public bool IsDeletedForUser(Guid userId) => _deletedFor.Exists(d => d.UserId == userId);
public virtual void Delete() => AddState(MessageState.IsDeleted);
public virtual void Edit(string newContent)
{
Content = newContent;
AddState(MessageState.IsEdited);
public virtual void Edit(string newContent)
{
Content = newContent;
AddState(MessageState.IsEdited);
}
public void DeleteForUser(Guid userId)
{
if (!_deletedFor.Exists(x => x.UserId == userId))
_deletedFor.Add(new DeletedMessage(Id, userId));
public void DeleteForUser(Guid userId)
{
if (!_deletedFor.Exists(x => x.UserId == userId))
_deletedFor.Add(new DeletedMessage(Id, userId));
}
public void MarkAsRead(Guid userId)
{
if (!_readByUsers.Contains(userId))
{
_readByUsers.Add(userId);
}
}
public bool IsReadBy(Guid userId) => _readByUsers.Contains(userId);
}
@@ -159,7 +159,7 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
(message as StoryMessage)?.StoryMediaType,
(message as MediaMessage)?.Media.Select(m => new MediaDto(m.Id, m.Type, m.Url, m.Filename, m.Size, m.Duration)).ToList() ?? new List<MediaDto>(),
sender != null ? new MessageSenderDto(sender.Id, sender.Username, sender.DisplayName, sender.Avatar) : new MessageSenderDto(message.SenderId, "unknown", "Unknown", null),
new List<ReadByDto>(), // ReadBy not implemented in this detailed view yet
message.ReadByUsers.Select(id => new ReadByDto(id)).ToList(),
reactions?.Select(r =>
{
senders.TryGetValue(r.UserId, out var ru);
@@ -1,7 +1,8 @@
using MediatR;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Shared.Kernel;
using MediatR;
namespace Knot.Modules.Conversations.Application.Messages.Read;
@@ -11,11 +12,13 @@ public sealed class ReadMessagesCommandHandler : ICommandHandler<ReadMessagesCom
{
private readonly IChatRepository _chatRepository;
private readonly IChatsUnitOfWork _unitOfWork;
private readonly IMessageRepository _messageRepository;
public ReadMessagesCommandHandler(IChatRepository chatRepository, IChatsUnitOfWork unitOfWork)
public ReadMessagesCommandHandler(IChatRepository chatRepository, IChatsUnitOfWork unitOfWork, IMessageRepository messageRepository)
{
_chatRepository = chatRepository;
_unitOfWork = unitOfWork;
_messageRepository = messageRepository;
}
public async Task<Result> Handle(ReadMessagesCommand request, CancellationToken cancellationToken)
@@ -28,6 +31,23 @@ public sealed class ReadMessagesCommandHandler : ICommandHandler<ReadMessagesCom
member.UpdateReadCursor(request.LastReadMessageId, request.LastReadSequenceId);
// Обновляем ReadByUsers для всех сообщений до LastReadSequenceId
var messages = await _messageRepository.GetChatMessagesAfterAsync(
request.ChatId,
0,
1000,
cancellationToken);
foreach (var message in messages)
{
if (message.SequenceId <= request.LastReadSequenceId &&
message.SenderId != request.UserId &&
!message.IsReadBy(request.UserId))
{
message.MarkAsRead(request.UserId);
}
}
await _unitOfWork.SaveChangesAsync(cancellationToken);
return Result.Success();
@@ -3,10 +3,10 @@ using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
@@ -41,7 +41,8 @@ internal sealed class SearchMessagesQueryHandler : IQueryHandler<SearchMessagesQ
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 => {
var result = messages.Select(message =>
{
var textMessage = message as TextMessage;
var mediaMessage = message as MediaMessage;
var storyMessage = message as StoryMessage;
@@ -66,7 +67,7 @@ internal sealed class SearchMessagesQueryHandler : IQueryHandler<SearchMessagesQ
mediaMessage?.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList() ?? new List<MediaDto>(),
senders.TryGetValue(message.SenderId, out var senderUser) ? new MessageSenderDto(senderUser.Id, senderUser.Username, senderUser.DisplayName, senderUser.Avatar) : new MessageSenderDto(message.SenderId, "unknown", "Unknown", null),
reactionsByMessage.TryGetValue(message.Id, out var mr) ? mr.Select(reaction => new SimpleReactionDto(reaction.UserId, reaction.Emoji)).ToList() : new List<SimpleReactionDto>(),
new List<ReadByDto>()
message.ReadByUsers.Select(id => new ReadByDto(id)).ToList()
);
}).ToList();
@@ -1,8 +1,8 @@
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Settings.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
namespace Knot.Modules.Conversations.Application.Messages.Send;
@@ -192,6 +192,7 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
var senderMember = chat.Members.First(m => m.UserId == request.SenderId);
senderMember.UpdateReadCursor(message.Id, message.SequenceId);
senderMember.UpdateDeliveredCursor(message.Id);
message.MarkAsRead(request.SenderId); // Отправитель всегда "прочитал" своё сообщение
// 5.
_messageRepository.Add(message);
@@ -1,26 +1,26 @@
using System.Collections.Concurrent;
using System.Security.Claims;
using Knot.Contracts.Auth.Application.Abstractions;
using Knot.Contracts.Auth.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Modules.Conversations.Application.Messages.Delete;
using Knot.Modules.Conversations.Application.Messages.Edit;
using Knot.Modules.Conversations.Application.Messages.Pin;
using Knot.Modules.Conversations.Application.Messages.React;
using Knot.Modules.Conversations.Application.Messages.Read;
using Knot.Modules.Conversations.Application.Messages.Send;
using Knot.Modules.Conversations.Application.Messages.Unpin;
using Knot.Modules.Conversations.Application.Messages.Vote;
using Knot.Shared.Kernel;
using MediatR;
using Microsoft.AspNetCore.Authorization;
using Microsoft.AspNetCore.SignalR;
using Microsoft.Extensions.Logging;
using Knot.Modules.Conversations.Application.Messages.Send;
using Knot.Modules.Conversations.Application.Messages.Read;
using Knot.Modules.Conversations.Application.Messages.Delete;
using Knot.Modules.Conversations.Application.Messages.React;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using Microsoft.Extensions.Caching.Memory;
using Knot.Contracts.Auth.Domain;
using Knot.Contracts.Auth.Application.Abstractions;
using Knot.Modules.Conversations.Application.Messages.Pin;
using Knot.Modules.Conversations.Application.Messages.Unpin;
using Knot.Modules.Conversations.Application.Messages.Vote;
using Knot.Modules.Conversations.Application.Messages.Edit;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Microsoft.Extensions.Logging;
namespace Knot.Modules.Conversations.Infrastructure.SignalR;
@@ -37,7 +37,7 @@ public sealed class ChatHub : Hub
public static int OnlineUsersCount => _userConnections.Count;
public static bool IsUserOnline(string userId) => _userConnections.ContainsKey(userId);
// userId → CallSession (one user can be in only one call at a time)
private static readonly ConcurrentDictionary<string, CallSession> _activeSessionsByUser = new();
// chatId → (startTime, callType)
@@ -53,12 +53,12 @@ public sealed class ChatHub : Hub
private readonly IUserDisplayNameProvider _userProvider;
public ChatHub(
ISender sender,
IUserContext userContext,
IChatRepository chatRepository,
IUserRepository userRepository,
ISender sender,
IUserContext userContext,
IChatRepository chatRepository,
IUserRepository userRepository,
IMessageRepository messageRepository,
ILogger<ChatHub> logger,
ILogger<ChatHub> logger,
IMemoryCache cache,
IUserDisplayNameProvider userProvider)
{
@@ -158,6 +158,7 @@ public sealed class ChatHub : Hub
await _sender.Send(command);
}
// Отправляем событие всем в чате о том, что пользователь прочитал сообщения
await Clients.Group(request.ChatId.ToString()).SendAsync("messages_read", new
{
ChatId = request.ChatId.ToString(),
@@ -261,7 +262,7 @@ public sealed class ChatHub : Hub
{
var senderInfo = await _userProvider.GetUsersInfoAsync(new[] { message.SenderId });
var dto = MessageMapper.MapToDto(message, senderInfo, Enumerable.Empty<MessageReaction>(), Enumerable.Empty<Guid>());
await Clients.Group(request.ChatId.ToString()).SendAsync("message_pinned", new
{
chatId = request.ChatId,
@@ -276,7 +277,7 @@ public sealed class ChatHub : Hub
{
var command = new UnpinMessageCommand(request.MessageId, request.ChatId, _userContext.UserId);
var result = await _sender.Send(command);
await Clients.Group(request.ChatId.ToString()).SendAsync("message_unpinned", new
{
chatId = request.ChatId,
@@ -328,7 +329,7 @@ public sealed class ChatHub : Hub
public async Task FriendAccepted(FriendSignalRequest request)
{
if (request == null || string.IsNullOrEmpty(request.FriendId)) return;
_logger.LogInformation("Signaling friend_request_accepted to {FriendId} from {UserId}", request.FriendId, _userContext.UserId);
await SendToUserAsync(request.FriendId, "friend_request_accepted", new { userId = _userContext.UserId });
}
@@ -350,8 +351,8 @@ public sealed class ChatHub : Hub
{
if (string.IsNullOrEmpty(targetUserId))
{
_logger.LogWarning("SendToUserAsync called with null or empty targetUserId");
return;
_logger.LogWarning("SendToUserAsync called with null or empty targetUserId");
return;
}
if (_userConnections.TryGetValue(targetUserId, out var connectionIds))
@@ -398,15 +399,15 @@ public sealed class ChatHub : Hub
// Track session for history
Guid? chatId = null;
if (Guid.TryParse(request.ChatId, out var parsedChatId)) chatId = parsedChatId;
if (!chatId.HasValue)
{
var userChats = await _chatRepository.GetUserChatsAsync(_userContext.UserId, Context.ConnectionAborted);
if (Guid.TryParse(request.TargetUserId, out var targetId))
{
var personalChat = userChats.FirstOrDefault(c => c.Type == ChatType.Personal && c.Members.Any(m => m.UserId == targetId));
if (personalChat != null) chatId = personalChat.Id;
}
var userChats = await _chatRepository.GetUserChatsAsync(_userContext.UserId, Context.ConnectionAborted);
if (Guid.TryParse(request.TargetUserId, out var targetId))
{
var personalChat = userChats.FirstOrDefault(c => c.Type == ChatType.Personal && c.Members.Any(m => m.UserId == targetId));
if (personalChat != null) chatId = personalChat.Id;
}
}
var session = new CallSession(chatId, _userContext.UserId, Guid.Parse(request.TargetUserId), request.CallType, DateTime.UtcNow);
@@ -461,10 +462,10 @@ public sealed class ChatHub : Hub
_activeSessionsByUser.TryRemove(request.TargetUserId, out _);
if (session.ChatId.HasValue)
{
int duration = session.IsAnswered && session.AnswerTime.HasValue
? (int)(DateTime.UtcNow - session.AnswerTime.Value).TotalSeconds
int duration = session.IsAnswered && session.AnswerTime.HasValue
? (int)(DateTime.UtcNow - session.AnswerTime.Value).TotalSeconds
: 0;
string status = session.IsAnswered ? "completed" : (_userContext.UserId == session.FromUserId ? "cancelled" : "missed");
await CreateCallMessage(session.ChatId.Value, session.FromUserId, session.CallType, status, duration);
}
@@ -573,7 +574,7 @@ public sealed class ChatHub : Hub
var userInfo = new ParticipantInfo(userId, username, displayName, avatar);
var participants = _groupCallParticipants.GetOrAdd(chatId, _ =>
var participants = _groupCallParticipants.GetOrAdd(chatId, _ =>
{
_activeGroupCalls[chatId] = (DateTime.UtcNow, request.CallType);
return new ConcurrentDictionary<string, ParticipantInfo>();
+24 -10
View File
@@ -18,7 +18,7 @@ public abstract class Message : AggregateRoot<Guid>
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;
@@ -26,10 +26,10 @@ public abstract class Message : AggregateRoot<Guid>
// ================== Опциональные метаданные (общего назначения) ==================
public Guid? ReplyToId { get; protected set; }
public Guid? ForwardedFromId { get; protected set; }
// ================== Флаги ==================
public MessageState State { get; protected set; }
// ================== Абстрактные / Виртуальные свойства ==================
public abstract string Type { get; }
public abstract string? Content { get; protected set; }
@@ -42,16 +42,20 @@ public abstract class Message : AggregateRoot<Guid>
protected List<DeletedMessage> _deletedFor = new();
public IReadOnlyCollection<DeletedMessage> DeletedFor => _deletedFor.AsReadOnly();
// ================== Прочитано ==================
protected List<Guid> _readByUsers = new();
public IReadOnlyCollection<Guid> ReadByUsers => _readByUsers.AsReadOnly();
// ================== Инфраструктурный конструктор EF ==================
protected Message() : base(Guid.Empty) { }
protected Message(
Guid id,
Guid chatId,
Guid senderId,
Guid? replyToId,
Guid? forwardedFromId,
DateTime createdAt,
Guid id,
Guid chatId,
Guid senderId,
Guid? replyToId,
Guid? forwardedFromId,
DateTime createdAt,
bool isImported) : base(id)
{
ChatId = chatId;
@@ -59,7 +63,7 @@ public abstract class Message : AggregateRoot<Guid>
ReplyToId = replyToId;
ForwardedFromId = forwardedFromId;
CreatedAt = createdAt;
if (isImported) AddState(MessageState.IsImported);
}
@@ -89,6 +93,16 @@ public abstract class Message : AggregateRoot<Guid>
_deletedFor.Add(new DeletedMessage(Id, userId));
}
}
public void MarkAsRead(Guid userId)
{
if (!_readByUsers.Contains(userId))
{
_readByUsers.Add(userId);
}
}
public bool IsReadBy(Guid userId) => _readByUsers.Contains(userId);
}
@@ -84,7 +84,7 @@ public sealed class MessageSentDomainEventHandler : INotificationHandler<Message
size = m.Size
}).ToList() ?? (object)Array.Empty<object>(),
sender = senderObj,
readBy = new List<object>(),
readBy = message.ReadByUsers.Select(id => new { id }).ToList(),
storyId = (message as StoryMessage)?.StoryId,
storyMediaUrl = (message as StoryMessage)?.StoryMediaUrl,
storyMediaType = (message as StoryMessage)?.StoryMediaType,