前端:

增加会话缓存
后端:
增加触发消息创建事件
增加消息创建事件触发后的处理函数
This commit is contained in:
2026-01-20 23:09:46 +08:00
parent be621e9ae2
commit bedcf97c9d
11 changed files with 158 additions and 88 deletions
@@ -30,58 +30,26 @@ namespace IM_API.Application.EventHandlers
_context = imContext;
_mapper = mapper;
}
/*
* 此方法有并发问题,当双方同时第一次发送消息时,
* 会出现同时创建的情况,其中一方会报错,
* 导致也未走到更新逻辑,会话丢失
*/
public async Task Handle(MessageCreatedEvent @event)
{
//此处仅处理私聊会话创建
if (@event.ChatType == ChatType.GROUP)
if(@event.ChatType == ChatType.PRIVATE)
{
return;
}
var conversation = await _context.Conversations.FirstOrDefaultAsync(
x => x.UserId == @event.MsgSenderId && x.TargetId == @event.MsgRecipientId
);
//如果首次发消息则创建双方会话
if (conversation is null)
{
Conversation senderCon = _mapper.Map<Conversation>(@event);
Conversation ReceptCon = _mapper.Map<Conversation>(@event);
ReceptCon.UserId = @event.MsgRecipientId;
ReceptCon.TargetId = @event.MsgSenderId;
ReceptCon.UnreadCount += 1;
ReceptCon.LastReadMessageId = null;
_context.Conversations.AddRange(senderCon,ReceptCon);
await _context.SaveChangesAsync();
}
else
{
Conversation senderCon = conversation;
Conversation? ReceptCon = await _context.Conversations.FirstOrDefaultAsync(
x => x.UserId == @event.MsgRecipientId && x.TargetId == @event.MsgSenderId);
if (ReceptCon is null)
Conversation? userAConversation = await _context.Conversations.FirstOrDefaultAsync(
x => x.UserId == @event.MsgSenderId && x.TargetId == @event.MsgRecipientId
);
Conversation? userBConversation = await _context.Conversations.FirstOrDefaultAsync(
x => x.UserId == @event.MsgRecipientId && x.TargetId == @event.MsgSenderId
);
if(userAConversation is null || userBConversation is null)
{
_logger.LogError("ConversationEventHandlerError:接收者会话对象缺失!Event:{Event}", JsonConvert.SerializeObject(@event));
throw new BaseException(CodeDefine.SYSTEM_ERROR);
_logger.LogError("消息事件更新会话信息失败:{@event}",@event);
}
//更新发送者conversation
senderCon.UnreadCount = 0;
senderCon.LastReadMessageId = @event.MessageId;
senderCon.LastMessage = @event.MessageContent;
senderCon.LastMessageTime = DateTime.Now;
//更新接收者conversation
ReceptCon.UnreadCount += 1;
ReceptCon.LastMessage = @event.MessageContent;
senderCon.LastMessageTime = DateTime.Now;
_context.Conversations.UpdateRange(senderCon, ReceptCon);
userAConversation.LastMessage = @event.MessageContent;
userAConversation.LastReadMessageId = @event.MessageId;
userBConversation.LastMessage = @event.MessageContent;
userBConversation.UnreadCount += 1;
_context.UpdateRange(userAConversation,userBConversation);
await _context.SaveChangesAsync();
}
}
}
@@ -1,6 +1,9 @@
using IM_API.Application.Interfaces;
using AutoMapper;
using IM_API.Application.Interfaces;
using IM_API.Domain.Events;
using IM_API.Dtos;
using IM_API.Hubs;
using IM_API.Models;
using IM_API.Tools;
using Microsoft.AspNetCore.SignalR;
@@ -9,15 +12,20 @@ namespace IM_API.Application.EventHandlers
public class SignalREventHandler : IEventHandler<MessageCreatedEvent>
{
private readonly IHubContext<ChatHub> _hub;
public SignalREventHandler(IHubContext<ChatHub> hub)
private readonly IMapper _mapper;
public SignalREventHandler(IHubContext<ChatHub> hub, IMapper mapper)
{
_hub = hub;
_mapper = mapper;
}
public async Task Handle(MessageCreatedEvent @event)
{
var streamKey = @event.StreamKey;
await _hub.Clients.Group(streamKey).SendAsync(SignalRMethodDefine.ReceiveMessage, @event);
if(@event.ChatType == Models.ChatType.PRIVATE)
{
MessageBaseDto messageBaseDto = _mapper.Map<MessageBaseDto>(_mapper.Map<Message>(@event));
await _hub.Clients.Users(@event.MsgRecipientId.ToString()).SendAsync("ReceiveMessage", messageBaseDto);
}
}
}
}
+9 -4
View File
@@ -1,4 +1,7 @@
using IM_API.Dtos;
using AutoMapper;
using IM_API.Application.Interfaces;
using IM_API.Domain.Events;
using IM_API.Dtos;
using IM_API.Interface.Services;
using IM_API.Models;
using IM_API.Tools;
@@ -11,10 +14,14 @@ namespace IM_API.Hubs
{
private IMessageSevice _messageService;
private readonly IConversationService _conversationService;
public ChatHub(IMessageSevice messageService, IConversationService conversationService)
private readonly IEventBus _eventBus;
private readonly IMapper _mapper;
public ChatHub(IMessageSevice messageService, IConversationService conversationService, IEventBus eventBus, IMapper mapper)
{
_messageService = messageService;
_conversationService = conversationService;
_eventBus = eventBus;
_mapper = mapper;
}
public async override Task OnConnectedAsync()
@@ -51,8 +58,6 @@ namespace IM_API.Hubs
{
msgInfo = await _messageService.SendGroupMessageAsync(int.Parse(userIdStr), dto.ReceiverId, dto);
}
await Clients.Users(dto.ReceiverId.ToString()).SendAsync("ReceiveMessage", msgInfo);
return;
}
}
+7 -1
View File
@@ -1,4 +1,6 @@
using AutoMapper;
using IM_API.Application.Interfaces;
using IM_API.Domain.Events;
using IM_API.Dtos;
using IM_API.Exceptions;
using IM_API.Interface.Services;
@@ -13,11 +15,13 @@ namespace IM_API.Services
private readonly ImContext _context;
private readonly ILogger<MessageService> _logger;
private readonly IMapper _mapper;
public MessageService(ImContext context, ILogger<MessageService> logger, IMapper mapper)
private readonly IEventBus _eventBus;
public MessageService(ImContext context, ILogger<MessageService> logger, IMapper mapper, IEventBus eventBus)
{
_context = context;
_logger = logger;
_mapper = mapper;
_eventBus = eventBus;
}
public async Task<List<MessageBaseDto>> GetMessagesAsync(int userId, int conversationId, int? msgId, int? pageSize, bool desc)
@@ -92,6 +96,7 @@ namespace IM_API.Services
message.StreamKey = StreamKeyBuilder.Group(groupId);
_context.Messages.Add(message);
await _context.SaveChangesAsync();
await _eventBus.PublishAsync(_mapper.Map<MessageCreatedEvent>(message));
return _mapper.Map<MessageBaseDto>(message);
}
@@ -107,6 +112,7 @@ namespace IM_API.Services
message.StreamKey = StreamKeyBuilder.Private(dto.SenderId, dto.ReceiverId);
_context.Messages.Add(message);
await _context.SaveChangesAsync();
await _eventBus.PublishAsync(_mapper.Map<MessageCreatedEvent>(message));
return _mapper.Map<MessageBaseDto>(message);
}
#endregion