using System.Text.RegularExpressions; using Knot.Contracts.Messaging.Application.Abstractions; using Knot.Contracts.Messaging.Domain; using Knot.Shared.Kernel; using Microsoft.Extensions.Configuration; using MongoDB.Bson; using MongoDB.Driver; using MongoDB.Driver.Linq; namespace Knot.Modules.Messaging.Infrastructure.Persistence; public sealed class MessageRepository : IMessageRepository { private readonly IMongoCollection _messages; private readonly IChatAccessProvider _chatAccessProvider; private readonly MediatR.IMediator _mediator; public MessageRepository(IMongoDatabase mongoDatabase, IChatAccessProvider chatAccessProvider, MediatR.IMediator mediator) { _messages = mongoDatabase.GetCollection("messages"); _chatAccessProvider = chatAccessProvider; _mediator = mediator; } public void Add(Message message) { _messages.InsertOne(message); // Publish domain events manualy for mongo entities var events = message.GetDomainEvents().ToList(); message.ClearDomainEvents(); // This runs synchronously or without waiting, better to run async but Add is void // In this implementation setting, fire and forget or wrap sync foreach (var domainEvent in events) { _mediator.Publish(domainEvent).GetAwaiter().GetResult(); } } public async Task GetByIdAsync(Guid id, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.Id, id); return await _messages.Find(filter).FirstOrDefaultAsync(cancellationToken); } public async Task> GetChatMessagesAsync(Guid chatId, int limit, int offset, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.ChatId, chatId); return await _messages.Find(filter) .SortByDescending(m => m.CreatedAt) .Skip(offset) .Limit(limit) .ToListAsync(cancellationToken); } public async Task GetLatestChatMessageAsync(Guid chatId, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.ChatId, chatId); return await _messages.Find(filter) .SortByDescending(m => m.CreatedAt) .FirstOrDefaultAsync(cancellationToken); } public async Task> GetPinnedMessagesAsync(Guid chatId, CancellationToken cancellationToken) { var builder = Builders.Filter; var filter = builder.And( builder.Eq(m => m.ChatId, chatId), builder.BitsAnySet(m => m.State, (long)MessageState.IsPinned) ); return await _messages.Find(filter) .SortByDescending(m => m.CreatedAt) .ToListAsync(cancellationToken); } public async Task> GetChatMessagesCursorAsync(Guid chatId, DateTime? cursor, long? sequenceId, int limit, CancellationToken cancellationToken) { var builder = Builders.Filter; var filter = builder.Eq(m => m.ChatId, chatId); if (sequenceId.HasValue) { filter &= builder.Lt(m => m.SequenceId, sequenceId.Value); } else if (cursor.HasValue) { filter &= builder.Lt(m => m.CreatedAt, cursor.Value); } return await _messages.Find(filter) .SortByDescending(m => m.SequenceId) .Limit(limit) .ToListAsync(cancellationToken); } public async Task> GetChatMessagesAroundAsync(Guid chatId, long sequenceId, int limit, CancellationToken cancellationToken) { var builder = Builders.Filter; // Target message var targetFilter = builder.And(builder.Eq(m => m.ChatId, chatId), builder.Eq(m => m.SequenceId, sequenceId)); var targetMsg = await _messages.Find(targetFilter).FirstOrDefaultAsync(cancellationToken); // Older messages var olderFilter = builder.And(builder.Eq(m => m.ChatId, chatId), builder.Lt(m => m.SequenceId, sequenceId)); var older = await _messages.Find(olderFilter) .SortByDescending(m => m.SequenceId) .Limit(limit / 2) .ToListAsync(cancellationToken); // Newer messages var newerFilter = builder.And(builder.Eq(m => m.ChatId, chatId), builder.Gt(m => m.SequenceId, sequenceId)); var newer = await _messages.Find(newerFilter) .SortBy(m => m.SequenceId) .Limit(limit / 2) .ToListAsync(cancellationToken); var result = new List(); result.AddRange(older); if (targetMsg != null) result.Add(targetMsg); result.AddRange(newer); return result.OrderBy(m => m.SequenceId).ToList(); } public async Task> SearchMessagesAsync(string query, Guid? chatId, Guid requestingUserId, CancellationToken cancellationToken) { // Not ideal for SQL/Mongo combination but keeping the signature var validChatIdsQuery = await _chatAccessProvider.GetValidChatIdsForUserAsync(requestingUserId, cancellationToken); var builder = Builders.Filter; var filter = builder.In(m => m.ChatId, validChatIdsQuery); if (chatId.HasValue) { filter &= builder.Eq(m => m.ChatId, chatId.Value); } var textFilter = Builders.Filter.Regex("Content", new BsonRegularExpression(Regex.Escape(query), "i")); filter &= textFilter; return await _messages.Find(filter) .SortByDescending(m => m.CreatedAt) .Limit(50) .ToListAsync(cancellationToken); } public async Task GetLastStoryMessageAsync(Guid chatId, Guid storyId, CancellationToken cancellationToken) { var filter = Builders.Filter.And( Builders.Filter.Eq(m => m.ChatId, chatId), Builders.Filter.Eq("_t", "StoryMessage"), Builders.Filter.Eq("StoryId", storyId) ); return await _messages.Find(filter) .SortByDescending(m => m.CreatedAt) .FirstOrDefaultAsync(cancellationToken); } public async Task UpdateAsync(Message message, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.Id, message.Id); await _messages.ReplaceOneAsync(filter, message, new ReplaceOptions { IsUpsert = true }, cancellationToken); // Publish domain events var events = message.GetDomainEvents().ToList(); message.ClearDomainEvents(); foreach (var domainEvent in events) { await _mediator.Publish(domainEvent, cancellationToken); } } public async Task> GetAllMessagesAsync(CancellationToken cancellationToken) { return await _messages.Find(_ => true).ToListAsync(cancellationToken); } public async Task DeleteChatMessagesAsync(Guid chatId, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.ChatId, chatId); await _messages.DeleteManyAsync(filter, cancellationToken); } public async Task DeleteUserMessagesAsync(Guid userId, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.SenderId, userId); await _messages.DeleteManyAsync(filter, cancellationToken); } public async Task RemoveAsync(Guid id, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.Id, id); await _messages.DeleteOneAsync(filter, cancellationToken); } public async Task DeleteAsync(Message message, CancellationToken cancellationToken) { var filter = Builders.Filter.Eq(m => m.Id, message.Id); await _messages.DeleteOneAsync(filter, cancellationToken); } }