Фиксы

This commit is contained in:
Халимов Рустам
2026-05-08 23:22:07 +03:00
parent a0b50b57b0
commit 00ed4b5959
12 changed files with 120 additions and 57 deletions

View File

@@ -41,6 +41,7 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
private readonly IChatsUnitOfWork _unitOfWork;
private readonly MediatR.IMediator _mediator;
private readonly IMessagesSettings _messagesSettings;
private readonly IIdempotencyStore _idempotencyStore;
private readonly ILogger<SendMessageCommandHandler> _logger;
public SendMessageCommandHandler(
@@ -49,6 +50,7 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
IChatsUnitOfWork unitOfWork,
MediatR.IMediator mediator,
IMessagesSettings messagesSettings,
IIdempotencyStore idempotencyStore,
ILogger<SendMessageCommandHandler> logger)
{
_chatRepository = chatRepository;
@@ -56,10 +58,25 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
_unitOfWork = unitOfWork;
_mediator = mediator;
_messagesSettings = messagesSettings;
_idempotencyStore = idempotencyStore;
_logger = logger;
}
public async Task<Result<Guid>> Handle(SendMessageCommand request, CancellationToken cancellationToken)
{
if (!string.IsNullOrWhiteSpace(request.IdempotencyKey))
{
var key = $"send_msg:{request.ChatId}:{request.IdempotencyKey}";
return await _idempotencyStore.GetOrCreateAsync(
key,
factory: ct => ExecuteAsync(request, ct),
cancellationToken: cancellationToken);
}
return await ExecuteAsync(request, cancellationToken);
}
private async Task<Result<Guid>> ExecuteAsync(SendMessageCommand request, CancellationToken cancellationToken)
{
// 1. Проверка существования чата
var chat = await _chatRepository.GetByIdAsync(request.ChatId, cancellationToken);
@@ -68,40 +85,30 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
return Result.Failure<Guid>(ChatErrors.ChatsNotFound);
}
// 2. <EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>, <20><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD> <20><> <20><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD> <20><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
// 2. Проверка, является ли отправитель участником чата
if (!chat.Members.Any(m => m.UserId == request.SenderId))
{
return Result.Failure<Guid>(ChatErrors.ChatsForbidden);
}
// 3. <EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD> <20><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
// 3. Создание сообщения
Message message;
if (request.Type == "story_reply" || request.Type == "story_reaction")
{
if (!_messagesSettings.Current.AllowMedia) return Result.Failure<Guid>(ChatErrors.MediaDisabled);
var parsedStoryMediaType = Enum.TryParse<MediaType>(request.StoryMediaType, true, out var sTypeEnum) ? sTypeEnum : MediaType.Image;
message = new StoryMessage(
Guid.NewGuid(),
request.ChatId,
request.SenderId,
request.StoryId ?? Guid.Empty,
request.StoryMediaUrl ?? string.Empty,
request.StoryMediaType,
request.Content,
request.ReplyToId,
request.ForwardedFromId,
DateTime.UtcNow,
false);
}
else if (request.Attachments != null && request.Attachments.Any())
@@ -111,26 +118,17 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
var firstAtt = request.Attachments.First();
var parsedType = Enum.TryParse<MediaType>(firstAtt.Type, true, out var mTypeEnum) ? mTypeEnum : MediaType.File;
message = new MediaMessage(
Guid.NewGuid(),
request.ChatId,
request.SenderId,
parsedType,
request.Content,
request.ReplyToId,
request.ForwardedFromId,
DateTime.UtcNow,
false);
foreach (var att in request.Attachments)
{
var pType = Enum.TryParse<MediaType>(att.Type, true, out var tEnum) ? tEnum : MediaType.File;
@@ -172,25 +170,17 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
{
message = new TextMessage(
Guid.NewGuid(),
request.ChatId,
request.SenderId,
request.Content ?? string.Empty,
request.ReplyToId,
request.Quote,
request.ForwardedFromId,
DateTime.UtcNow,
false);
}
// 4. <EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD> <20><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD> High-Water Mark
// 4. Обновление High-Water Mark
chat.IncrementSequenceId();
message.SetSequenceId(chat.LastMessageSequenceId);
@@ -204,13 +194,9 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
await _mediator.Publish(new MessageSentDomainEvent(
message.Id,
message.ChatId,
message.SenderId,
message.Content),
cancellationToken);
return Result.Success(message.Id);

View File

@@ -2,7 +2,9 @@ using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Infrastructure.Persistence;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Modules.Conversations.Infrastructure.Idempotency;
using Knot.Modules.Conversations.Infrastructure.Persistence;
using Knot.Modules.Conversations.Infrastructure.Persistence.Mongo;
using Knot.Modules.Conversations.Infrastructure.Services;
using Knot.Shared.Kernel;
using Microsoft.EntityFrameworkCore;
@@ -42,6 +44,9 @@ public static class DependencyInjection
services.AddScoped<Knot.Contracts.Conversations.Abstractions.IUserStatusService, UserStatusService>();
services.AddScoped<Knot.Contracts.Conversations.Abstractions.IUserDeleterService, UserDeleterService>();
services.AddMemoryCache();
services.AddSingleton<Knot.Contracts.Conversations.Application.Abstractions.IIdempotencyStore, MemoryCacheIdempotencyStore>();
return services;
}
}

View File

@@ -0,0 +1,54 @@
using Knot.Contracts.Conversations.Application.Abstractions;
using Microsoft.Extensions.Caching.Memory;
namespace Knot.Modules.Conversations.Infrastructure.Idempotency;
/// <summary>
/// Реализация хранилища идемпотентности на основе IMemoryCache.
/// Использует семафор для предотвращения race condition при одновременных запросах с одинаковым ключом.
/// </summary>
public sealed class MemoryCacheIdempotencyStore : IIdempotencyStore
{
private readonly IMemoryCache _cache;
private readonly SemaphoreSlim _semaphore = new(1, 1);
public MemoryCacheIdempotencyStore(IMemoryCache cache)
{
_cache = cache;
}
public async Task<T> GetOrCreateAsync<T>(
string key,
Func<CancellationToken, Task<T>> factory,
TimeSpan? expiration = null,
CancellationToken cancellationToken = default)
{
if (_cache.TryGetValue(key, out T? cachedValue) && cachedValue is not null)
{
return cachedValue;
}
await _semaphore.WaitAsync(cancellationToken);
try
{
// Double-check после получения блокировки
if (_cache.TryGetValue(key, out cachedValue) && cachedValue is not null)
{
return cachedValue;
}
var value = await factory(cancellationToken);
var options = new MemoryCacheEntryOptions()
.SetAbsoluteExpiration(expiration ?? TimeSpan.FromHours(24))
.SetPriority(CacheItemPriority.Normal);
_cache.Set(key, value, options);
return value;
}
finally
{
_semaphore.Release();
}
}
}

View File

@@ -29,7 +29,6 @@
<PackageReference Include="Npgsql.EntityFrameworkCore.PostgreSQL" Version="10.0.1" />
<PackageReference Include="MongoDB.Driver" Version="3.2.0" />
<PackageReference Include="SixLabors.ImageSharp" Version="3.1.12" />
<PackageReference Include="DKNet.AspCore.Idempotency" Version="1.0.0" />
</ItemGroup>
<ItemGroup>
@@ -38,6 +37,7 @@
<ItemGroup>
<InternalsVisibleTo Include="DynamicProxyGenAssembly2" />
<InternalsVisibleTo Include="Knot.Modules.Conversations.UnitTests" />
</ItemGroup>
</Project>