From 0c5a15d7af003862b34ab3b9691f915d896b5227 Mon Sep 17 00:00:00 2001 From: Stepan Grankin Date: Fri, 7 Aug 2026 22:19:16 +0300 Subject: [PATCH 1/6] =?UTF-8?q?feat(services):=20=D0=BE=D1=87=D0=B5=D1=80?= =?UTF-8?q?=D0=B5=D0=B4=D1=8C=20=D1=84=D0=BE=D0=BD=D0=BE=D0=B2=D1=8B=D1=85?= =?UTF-8?q?=20=D1=80=D0=B0=D0=B1=D0=BE=D1=82=20=D0=BD=D0=B0=20Channel?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit BackgroundTaskQueue поверх ограниченного Channel плюс BackgroundTaskProcessor как BackgroundService. Отличия от прежнего fire-and-forget: - очередь ограничена (1000), переполнение не блокирует поток запроса, а отбрасывает работу с предупреждением в лог - параллелизм обработчика ограничен четырьмя работами - при остановке приложения очередь закрывается на запись, а уже принятые работы дочитываются и выполняются, а не теряются вместе с процессом Вызовы Forget() пока не тронуты — переключение отдельным коммитом. --- src/Directory.Packages.props | 1 + .../DI/InternalServicesRegistration.cs | 5 + .../BackgroundTaskProcessorTests.cs | 163 ++++++++++++++++++ .../BackgroundTasks/BackgroundTask.cs | 18 ++ .../BackgroundTaskProcessor.cs | 56 ++++++ .../BackgroundTasks/BackgroundTaskQueue.cs | 70 ++++++++ .../BackgroundTasks/IBackgroundTaskQueue.cs | 21 +++ .../FillInTheTextBot.Services.csproj | 1 + 8 files changed, 335 insertions(+) create mode 100644 src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs create mode 100644 src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTask.cs create mode 100644 src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs create mode 100644 src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs create mode 100644 src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskQueue.cs diff --git a/src/Directory.Packages.props b/src/Directory.Packages.props index 6dcb0980..8419fa78 100644 --- a/src/Directory.Packages.props +++ b/src/Directory.Packages.props @@ -6,6 +6,7 @@ + diff --git a/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs b/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs index 46b71334..03772694 100644 --- a/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs +++ b/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs @@ -1,4 +1,5 @@ using FillInTheTextBot.Services; +using FillInTheTextBot.Services.BackgroundTasks; using Microsoft.Extensions.DependencyInjection; namespace FillInTheTextBot.Api.DI @@ -9,6 +10,10 @@ internal static void AddInternalServices(this IServiceCollection services) { services.AddTransient(); services.AddScoped(); + + services.AddSingleton(); + services.AddSingleton(provider => provider.GetRequiredService()); + services.AddHostedService(); } } } diff --git a/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs new file mode 100644 index 00000000..976cf100 --- /dev/null +++ b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs @@ -0,0 +1,163 @@ +using System; +using System.Collections.Concurrent; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using FillInTheTextBot.Services.BackgroundTasks; +using Microsoft.Extensions.Logging.Abstractions; +using NUnit.Framework; +using NUnit.Framework.Legacy; + +namespace FillInTheTextBot.Services.Tests.BackgroundTasks +{ + [TestFixture] + public class BackgroundTaskProcessorTests + { + private BackgroundTaskQueue _queue; + private BackgroundTaskProcessor _target; + + [SetUp] + public void InitTest() + { + _queue = new BackgroundTaskQueue(NullLogger.Instance); + _target = new BackgroundTaskProcessor(_queue, NullLogger.Instance); + } + + [TearDown] + public async Task CleanUp() + { + await _target.StopAsync(CancellationToken.None); + _target.Dispose(); + } + + [Test] + public async Task Enqueue_Work_Executed() + { + var executed = new TaskCompletionSource(); + + await _target.StartAsync(CancellationToken.None); + + var accepted = _queue.Enqueue("проверка", () => + { + executed.TrySetResult(true); + return Task.CompletedTask; + }); + + + var completed = await Task.WhenAny(executed.Task, Task.Delay(5000)); + + + ClassicAssert.True(accepted); + ClassicAssert.AreSame(executed.Task, completed, "Работа из очереди должна быть выполнена"); + } + + [Test] + public async Task Enqueue_ManyWorks_AllExecuted() + { + const int count = 50; + + var executed = new ConcurrentBag(); + + await _target.StartAsync(CancellationToken.None); + + foreach (var i in Enumerable.Range(0, count)) + { + var number = i; + + _queue.Enqueue($"работа-{number}", () => + { + executed.Add(number); + return Task.CompletedTask; + }); + } + + + await WaitForAsync(() => executed.Count == count); + + + ClassicAssert.AreEqual(count, executed.Count); + CollectionAssert.AreEquivalent(Enumerable.Range(0, count), executed); + } + + [Test] + public async Task Enqueue_FailingWork_QueueKeepsWorking() + { + var executed = new TaskCompletionSource(); + + await _target.StartAsync(CancellationToken.None); + + _queue.Enqueue("падающая", () => throw new InvalidOperationException("ошибка внутри работы")); + + _queue.Enqueue("следующая", () => + { + executed.TrySetResult(true); + return Task.CompletedTask; + }); + + + var completed = await Task.WhenAny(executed.Task, Task.Delay(5000)); + + + ClassicAssert.AreSame(executed.Task, completed, "Исключение в одной работе не должно останавливать обработчик"); + } + + [Test] + public async Task StopAsync_PendingWork_Executed() + { + var executed = new TaskCompletionSource(); + + await _target.StartAsync(CancellationToken.None); + + _queue.Enqueue("до остановки", () => + { + executed.TrySetResult(true); + return Task.CompletedTask; + }); + + + await _target.StopAsync(CancellationToken.None); + + + var completed = await Task.WhenAny(executed.Task, Task.Delay(5000)); + + ClassicAssert.AreSame(executed.Task, completed, "Принятые работы должны успеть выполниться при остановке"); + } + + [Test] + public void Enqueue_QueueIsFull_TaskDropped() + { + // Обработчик не запущен, поэтому очередь только наполняется + foreach (var i in Enumerable.Range(0, BackgroundTaskQueue.Capacity)) + { + var accepted = _queue.Enqueue($"работа-{i}", () => Task.CompletedTask); + + ClassicAssert.True(accepted, $"Работа {i} должна помещаться в очередь"); + } + + + var overflow = _queue.Enqueue("лишняя", () => Task.CompletedTask); + + + ClassicAssert.False(overflow, "Переполнение очереди не должно блокировать вызывающий поток"); + } + + [Test] + public void Enqueue_NullWork_NotAccepted() + { + var accepted = _queue.Enqueue("пустая", null); + + ClassicAssert.False(accepted); + } + + private static async Task WaitForAsync(Func condition, int timeoutMilliseconds = 5000) + { + var waited = 0; + + while (!condition() && waited < timeoutMilliseconds) + { + await Task.Delay(25); + waited += 25; + } + } + } +} diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTask.cs b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTask.cs new file mode 100644 index 00000000..e1e52064 --- /dev/null +++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTask.cs @@ -0,0 +1,18 @@ +using System; +using System.Threading.Tasks; + +namespace FillInTheTextBot.Services.BackgroundTasks +{ + public sealed class BackgroundTask + { + public BackgroundTask(string name, Func work) + { + Name = name; + Work = work; + } + + public string Name { get; } + + public Func Work { get; } + } +} diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs new file mode 100644 index 00000000..702ca31b --- /dev/null +++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs @@ -0,0 +1,56 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; + +namespace FillInTheTextBot.Services.BackgroundTasks +{ + public sealed class BackgroundTaskProcessor : BackgroundService + { + /// + /// Работы в очереди независимы, поэтому выполняются параллельно — как это было + /// с прежним fire-and-forget. Ограничение не даёт при всплеске открыть + /// неограниченное число обращений к Redis и Dialogflow. + /// + public const int MaxDegreeOfParallelism = 4; + + private readonly BackgroundTaskQueue _queue; + private readonly ILogger _log; + + public BackgroundTaskProcessor(BackgroundTaskQueue queue, ILogger log) + { + _queue = queue; + _log = log; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + var options = new ParallelOptions { MaxDegreeOfParallelism = MaxDegreeOfParallelism }; + + // Чтение намеренно не отменяется по stoppingToken: цикл заканчивается, когда + // очередь закрыта на запись и разобрана. Так при остановке приложения уже + // принятые работы доводятся до конца, а не теряются + await Parallel.ForEachAsync(_queue.ReadAllAsync(), options, (task, _) => ExecuteTaskAsync(task)); + } + + public override async Task StopAsync(CancellationToken cancellationToken) + { + _queue.Complete(); + + await base.StopAsync(cancellationToken); + } + + private async ValueTask ExecuteTaskAsync(BackgroundTask task) + { + try + { + await task.Work().ConfigureAwait(false); + } + catch (Exception e) + { + _log.LogError(e, "Error while executing background task '{TaskName}'", task.Name); + } + } + } +} diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs new file mode 100644 index 00000000..e00acb1b --- /dev/null +++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs @@ -0,0 +1,70 @@ +using System; +using System.Collections.Generic; +using System.Threading.Channels; +using System.Threading.Tasks; +using Microsoft.Extensions.Logging; + +namespace FillInTheTextBot.Services.BackgroundTasks +{ + public sealed class BackgroundTaskQueue : IBackgroundTaskQueue + { + /// + /// Ёмкость подобрана с запасом на всплеск запросов: при штатной нагрузке очередь + /// разбирается быстрее, чем наполняется, а при недоступности внешнего сервиса + /// ограничение не даёт очереди съесть память. + /// + public const int Capacity = 1000; + + private readonly Channel _channel; + private readonly ILogger _log; + + public BackgroundTaskQueue(ILogger log) + { + _log = log; + + // Режим Wait выбран ради TryWrite: на заполненной очереди он возвращает false, + // не блокируя вызывающий поток, и переполнение видно в логе. Режим DropWrite + // молча отбрасывал бы работу и возвращал true + var options = new BoundedChannelOptions(Capacity) + { + FullMode = BoundedChannelFullMode.Wait, + SingleReader = false, + SingleWriter = false + }; + + _channel = Channel.CreateBounded(options); + } + + public bool Enqueue(string name, Func work) + { + if (work is null) + { + return false; + } + + var task = new BackgroundTask(name, work); + + var written = _channel.Writer.TryWrite(task); + + if (!written) + { + _log.LogWarning("Background task queue is full, task '{TaskName}' is dropped", name); + } + + return written; + } + + public IAsyncEnumerable ReadAllAsync() + { + return _channel.Reader.ReadAllAsync(); + } + + /// + /// Закрывает очередь на запись. Уже принятые работы остаются доступны для чтения. + /// + public void Complete() + { + _channel.Writer.TryComplete(); + } + } +} diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskQueue.cs b/src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskQueue.cs new file mode 100644 index 00000000..f2cf53bc --- /dev/null +++ b/src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskQueue.cs @@ -0,0 +1,21 @@ +using System; +using System.Threading.Tasks; + +namespace FillInTheTextBot.Services.BackgroundTasks +{ + /// + /// Очередь работ, которые не нужно дожидаться в рамках обработки запроса + /// (запись в кэш, установка контекста Dialogflow). + /// + public interface IBackgroundTaskQueue + { + /// + /// Ставит работу в очередь. Не блокирует вызывающий поток: если очередь переполнена, + /// работа отбрасывается — задержать ответ пользователю хуже, чем потерять запись. + /// + /// Имя работы, попадает в лог при ошибке или отбрасывании. + /// Работа. + /// true, если работа принята в очередь. + bool Enqueue(string name, Func work); + } +} diff --git a/src/FillInTheTextBot.Services/FillInTheTextBot.Services.csproj b/src/FillInTheTextBot.Services/FillInTheTextBot.Services.csproj index bb8a7d29..a3c38f5a 100644 --- a/src/FillInTheTextBot.Services/FillInTheTextBot.Services.csproj +++ b/src/FillInTheTextBot.Services/FillInTheTextBot.Services.csproj @@ -9,6 +9,7 @@ + From cb8a701c1942fc1cc9017091d13fc7660cbd6e44 Mon Sep 17 00:00:00 2001 From: Stepan Grankin Date: Fri, 7 Aug 2026 23:16:48 +0300 Subject: [PATCH 2/6] =?UTF-8?q?fix(services):=20=D0=BD=D0=B0=D0=B4=D1=91?= =?UTF-8?q?=D0=B6=D0=BD=D1=8B=D0=B9=20=D1=80=D0=B0=D0=B7=D0=B1=D0=BE=D1=80?= =?UTF-8?q?=20=D0=BE=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D0=B8=20=D1=84=D0=BE?= =?UTF-8?q?=D0=BD=D0=BE=D0=B2=D1=8B=D1=85=20=D1=80=D0=B0=D0=B1=D0=BE=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Две проблемы обработчика, найденные замерами. Первая — потеря работ при остановке. BackgroundService запускает ExecuteAsync отложенно: если приложение останавливается вскоре после старта, задача отменяется до входа в тело метода, и очередь никто не разбирает. Из 300 прогонов падало 263. Разбор вынесен в отдельный метод, который вызывается и из ExecuteAsync, и из StopAsync после базовой остановки — стало 0 из 300. Вторая — потеря параллельности. Прежний fire-and-forget выполнял работы одновременно, и медленная не задерживала остальные. Очередь с одним потребителем это ломала: работа, поставленная за вызовом Dialogflow, ждала его 200 мс. Теперь канал читают несколько независимых потребителей, их число ограничено четырьмя. Замер после правки: 0 мс. Parallel.ForEachAsync для этого не подходит: при остановке его задача завершается как отменённая, не дочитав очередь. Добавлен тест-сторож на параллельность: при одном потребителе он краснеет. --- .../BackgroundTaskProcessorTests.cs | 35 +++++++++++++ .../BackgroundTaskProcessor.cs | 51 +++++++++++++++---- 2 files changed, 75 insertions(+), 11 deletions(-) diff --git a/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs index 976cf100..b928c935 100644 --- a/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs +++ b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs @@ -101,6 +101,41 @@ public async Task Enqueue_FailingWork_QueueKeepsWorking() ClassicAssert.AreSame(executed.Task, completed, "Исключение в одной работе не должно останавливать обработчик"); } + [Test] + public async Task Enqueue_WorkBehindSlowOne_NotBlocked() + { + // Прежний fire-and-forget выполнял работы одновременно, и медленная не + // задерживала остальные. Это свойство должно сохраняться + var slowStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var slowFinish = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var fastExecuted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + await _target.StartAsync(CancellationToken.None); + + _queue.Enqueue("медленная", () => + { + slowStarted.TrySetResult(true); + return slowFinish.Task; + }); + + await slowStarted.Task; + + _queue.Enqueue("быстрая", () => + { + fastExecuted.TrySetResult(true); + return Task.CompletedTask; + }); + + + var completed = await Task.WhenAny(fastExecuted.Task, Task.Delay(5000)); + + + slowFinish.TrySetResult(true); + + ClassicAssert.AreSame(fastExecuted.Task, completed, + "Быстрая работа не должна ждать завершения медленной"); + } + [Test] public async Task StopAsync_PendingWork_Executed() { diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs index 702ca31b..ae225062 100644 --- a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs +++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs @@ -1,4 +1,6 @@ using System; +using System.Collections.Generic; +using System.Linq; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.Hosting; @@ -9,9 +11,10 @@ namespace FillInTheTextBot.Services.BackgroundTasks public sealed class BackgroundTaskProcessor : BackgroundService { /// - /// Работы в очереди независимы, поэтому выполняются параллельно — как это было - /// с прежним fire-and-forget. Ограничение не даёт при всплеске открыть - /// неограниченное число обращений к Redis и Dialogflow. + /// Сколько работ выполняется одновременно. Прежний fire-and-forget запускал их + /// без ограничений, и медленная работа не задерживала остальные — это свойство + /// нужно сохранить. Ограничение не даёт при всплеске открыть неограниченное + /// число обращений к Redis и Dialogflow. /// public const int MaxDegreeOfParallelism = 4; @@ -24,14 +27,9 @@ public BackgroundTaskProcessor(BackgroundTaskQueue queue, ILogger ExecuteTaskAsync(task)); + return DrainAsync(); } public override async Task StopAsync(CancellationToken cancellationToken) @@ -39,9 +37,40 @@ public override async Task StopAsync(CancellationToken cancellationToken) _queue.Complete(); await base.StopAsync(cancellationToken); + + // BackgroundService запускает ExecuteAsync отложенно, и при остановке вскоре + // после старта задача отменяется до входа в тело метода — тогда очередь никто + // не разобрал. Поэтому остаток добирается здесь, независимо от того, + // успел ли стартовать основной цикл. + await DrainAsync(); + } + + private Task DrainAsync() + { + // Несколько независимых потребителей одного канала. Parallel.ForEachAsync здесь + // не подходит: при остановке его задача завершается как отменённая, не дочитав + // очередь, и принятые работы теряются. + var consumers = Enumerable + .Range(0, MaxDegreeOfParallelism) + .Select(_ => ConsumeAsync()); + + return Task.WhenAll(consumers); + } + + /// + /// Читает очередь, пока она не закрыта на запись и не разобрана. Отмена по + /// stoppingToken намеренно не используется: при остановке приложения уже принятые + /// работы должны быть доведены до конца, а не потеряны вместе с процессом. + /// + private async Task ConsumeAsync() + { + await foreach (var task in _queue.ReadAllAsync()) + { + await ExecuteTaskAsync(task); + } } - private async ValueTask ExecuteTaskAsync(BackgroundTask task) + private async Task ExecuteTaskAsync(BackgroundTask task) { try { From a0af10d70a9172a435007797b3ee99e3ca0dde48 Mon Sep 17 00:00:00 2001 From: Stepan Grankin Date: Fri, 7 Aug 2026 23:16:48 +0300 Subject: [PATCH 3/6] =?UTF-8?q?refactor:=20=D0=BF=D0=B5=D1=80=D0=B5=D0=B2?= =?UTF-8?q?=D0=BE=D0=B4=20fire-and-forget=20=D0=BD=D0=B0=20=D0=BE=D1=87?= =?UTF-8?q?=D0=B5=D1=80=D0=B5=D0=B4=D1=8C=20=D1=84=D0=BE=D0=BD=D0=BE=D0=B2?= =?UTF-8?q?=D1=8B=D1=85=20=D1=80=D0=B0=D0=B1=D0=BE=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Три места, где результат намеренно не дожидались, переведены с TasksExtensions.Forget() на IBackgroundTaskQueue: - ConversationService — установка контекста savedText - SberService — запись идентификатора сессии в Redis - MarusiaService — отметка о пользователе в Redis TasksExtensions удалён. Прежний Forget() запускал задачу без ограничений по количеству и терял её при остановке процесса; теперь работы идут через ограниченную очередь и дочитываются при штатной остановке. Постановка в очередь не блокирует поток запроса вообще — это запись в канал. Forget() был чуть хуже: асинхронный метод стартовал синхронно на потоке запроса до первого await. В тестах подменена сама очередь: она выполняет работу сразу, поэтому ожидание с опросом больше не нужно. Проверки остались прежними — SetContextAsync и AddAsync вызываются с теми же аргументами. --- .../MarusiaServiceTests.cs | 15 +++++++- .../MarusiaService.cs | 12 ++++-- .../SberServiceTests.cs | 15 +++++++- .../SberService.cs | 12 ++++-- .../ConversationServiceTests.cs | 16 +++++++- .../ConversationService.cs | 8 +++- .../Extensions/TasksExtensions.cs | 37 ------------------- 7 files changed, 67 insertions(+), 48 deletions(-) delete mode 100644 src/FillInTheTextBot.Services/Extensions/TasksExtensions.cs diff --git a/src/FillInTheTextBot.Messengers.Marusia.Tests/MarusiaServiceTests.cs b/src/FillInTheTextBot.Messengers.Marusia.Tests/MarusiaServiceTests.cs index 3882d6b8..65394281 100644 --- a/src/FillInTheTextBot.Messengers.Marusia.Tests/MarusiaServiceTests.cs +++ b/src/FillInTheTextBot.Messengers.Marusia.Tests/MarusiaServiceTests.cs @@ -3,6 +3,7 @@ using System.Threading.Tasks; using AutoFixture; using FillInTheTextBot.Services; +using FillInTheTextBot.Services.BackgroundTasks; using GranSteL.Helpers.Redis; using MailRu.Marusia.Models; using MailRu.Marusia.Models.Input; @@ -23,6 +24,7 @@ public class MarusiaServiceTests { private Mock _conversationService; private Mock _cache; + private Mock _backgroundTasks; private MarusiaService _target; @@ -34,7 +36,18 @@ public void InitTest() _conversationService = new Mock(); _cache = new Mock(); - _target = new MarusiaService(Mock.Of>(), _conversationService.Object, _cache.Object); + _backgroundTasks = new Mock(); + + // Фоновые работы выполняются сразу, чтобы тест проверял результат, а не гонку + _backgroundTasks + .Setup(q => q.Enqueue(It.IsAny(), It.IsAny>())) + .Returns((string _, Func work) => + { + work().GetAwaiter().GetResult(); + return true; + }); + + _target = new MarusiaService(Mock.Of>(), _conversationService.Object, _cache.Object, _backgroundTasks.Object); _fixture = new Fixture { OmitAutoProperties = true }; } diff --git a/src/FillInTheTextBot.Messengers.Marusia/MarusiaService.cs b/src/FillInTheTextBot.Messengers.Marusia/MarusiaService.cs index 8edc5aee..5bb305e7 100644 --- a/src/FillInTheTextBot.Messengers.Marusia/MarusiaService.cs +++ b/src/FillInTheTextBot.Messengers.Marusia/MarusiaService.cs @@ -1,7 +1,7 @@ using System; using System.Threading.Tasks; using FillInTheTextBot.Services; -using FillInTheTextBot.Services.Extensions; +using FillInTheTextBot.Services.BackgroundTasks; using GranSteL.Helpers.Redis; using MailRu.Marusia.Models; using MailRu.Marusia.Models.Input; @@ -15,13 +15,16 @@ public class MarusiaService : MessengerService, IMarusi private const string PongResponse = "pong"; private readonly IRedisCacheService _cache; + private readonly IBackgroundTaskQueue _backgroundTasks; public MarusiaService( ILogger log, IConversationService conversationService, - IRedisCacheService cache) : base(log, conversationService) + IRedisCacheService cache, + IBackgroundTaskQueue backgroundTasks) : base(log, conversationService) { _cache = cache; + _backgroundTasks = backgroundTasks; } protected override Models.Request Before(InputModel input) @@ -67,7 +70,10 @@ protected override Task AfterAsync(InputModel input, Models.Respons output.AddToSessionState(Models.Response.ScopeStorageKey, response.ScopeKey); - _cache.AddAsync($"marusia:{input.Session?.UserId}", string.Empty, TimeSpan.FromDays(14)).Forget(); + var cacheKey = $"marusia:{input.Session?.UserId}"; + + _backgroundTasks.Enqueue("marusia user mark", + () => _cache.AddAsync(cacheKey, string.Empty, TimeSpan.FromDays(14))); return Task.FromResult(output); } diff --git a/src/FillInTheTextBot.Messengers.Sber.Tests/SberServiceTests.cs b/src/FillInTheTextBot.Messengers.Sber.Tests/SberServiceTests.cs index d34cf72d..d6328369 100644 --- a/src/FillInTheTextBot.Messengers.Sber.Tests/SberServiceTests.cs +++ b/src/FillInTheTextBot.Messengers.Sber.Tests/SberServiceTests.cs @@ -2,6 +2,7 @@ using System.Linq; using System.Threading.Tasks; using AutoFixture; +using FillInTheTextBot.Services.BackgroundTasks; using GranSteL.Helpers.Redis; using Microsoft.Extensions.Logging; using Moq; @@ -21,6 +22,7 @@ public class SberServiceTests { private Mock _conversationService; private Mock _cache; + private Mock _backgroundTasks; private SberService _target; @@ -32,7 +34,18 @@ public void InitTest() _conversationService = new Mock(); _cache = new Mock(); - _target = new SberService(Mock.Of>(), _conversationService.Object, _cache.Object); + _backgroundTasks = new Mock(); + + // Фоновые работы выполняются сразу, чтобы тест проверял результат, а не гонку + _backgroundTasks + .Setup(q => q.Enqueue(It.IsAny(), It.IsAny>())) + .Returns((string _, Func work) => + { + work().GetAwaiter().GetResult(); + return true; + }); + + _target = new SberService(Mock.Of>(), _conversationService.Object, _cache.Object, _backgroundTasks.Object); _fixture = new Fixture { OmitAutoProperties = true }; } diff --git a/src/FillInTheTextBot.Messengers.Sber/SberService.cs b/src/FillInTheTextBot.Messengers.Sber/SberService.cs index c3ec2f79..87798055 100644 --- a/src/FillInTheTextBot.Messengers.Sber/SberService.cs +++ b/src/FillInTheTextBot.Messengers.Sber/SberService.cs @@ -2,7 +2,7 @@ using System.Collections.Generic; using System.Threading.Tasks; using FillInTheTextBot.Services; -using FillInTheTextBot.Services.Extensions; +using FillInTheTextBot.Services.BackgroundTasks; using GranSteL.Helpers.Redis; using Microsoft.Extensions.Logging; using Sber.SmartApp.Models; @@ -12,13 +12,16 @@ namespace FillInTheTextBot.Messengers.Sber public class SberService : MessengerService, ISberService { private readonly IRedisCacheService _cache; + private readonly IBackgroundTaskQueue _backgroundTasks; public SberService( ILogger log, IConversationService conversationService, - IRedisCacheService cache) : base(log, conversationService) + IRedisCacheService cache, + IBackgroundTaskQueue backgroundTasks) : base(log, conversationService) { _cache = cache; + _backgroundTasks = backgroundTasks; } protected override Models.Request Before(Request input) @@ -76,7 +79,10 @@ private string TryGetSessionIdAsync(bool? newSession, string userHash) { sessionId = Guid.NewGuid().ToString("N"); - _cache.AddAsync(cacheKey, sessionId, TimeSpan.FromMinutes(5)).Forget(); + var newSessionId = sessionId; + + _backgroundTasks.Enqueue("sber session", + () => _cache.AddAsync(cacheKey, newSessionId, TimeSpan.FromMinutes(5))); } return sessionId; diff --git a/src/FillInTheTextBot.Services.Tests/ConversationServiceTests.cs b/src/FillInTheTextBot.Services.Tests/ConversationServiceTests.cs index c528a908..956cfed8 100644 --- a/src/FillInTheTextBot.Services.Tests/ConversationServiceTests.cs +++ b/src/FillInTheTextBot.Services.Tests/ConversationServiceTests.cs @@ -2,7 +2,9 @@ using System.Linq; using System.Threading.Tasks; using AutoFixture; +using System; using FillInTheTextBot.Models; +using FillInTheTextBot.Services.BackgroundTasks; using FillInTheTextBot.Services.Configuration; using GranSteL.Helpers.Redis; using Moq; @@ -20,6 +22,7 @@ public class ConversationServiceTests { private Mock _dialogflowService; private Mock _cache; + private Mock _backgroundTasks; private ConversationConfiguration _configuration; @@ -33,9 +36,20 @@ public void InitTest() _dialogflowService = new Mock(); _cache = new Mock(); + _backgroundTasks = new Mock(); + + // Фоновые работы выполняются сразу, чтобы тест проверял результат, а не гонку + _backgroundTasks + .Setup(q => q.Enqueue(It.IsAny(), It.IsAny>())) + .Returns((string _, Func work) => + { + work().GetAwaiter().GetResult(); + return true; + }); + _configuration = new ConversationConfiguration(); - _target = new ConversationService(_configuration, _dialogflowService.Object, _cache.Object); + _target = new ConversationService(_configuration, _dialogflowService.Object, _cache.Object, _backgroundTasks.Object); _fixture = new Fixture { OmitAutoProperties = true }; } diff --git a/src/FillInTheTextBot.Services/ConversationService.cs b/src/FillInTheTextBot.Services/ConversationService.cs index 21611656..7afa1911 100644 --- a/src/FillInTheTextBot.Services/ConversationService.cs +++ b/src/FillInTheTextBot.Services/ConversationService.cs @@ -3,6 +3,7 @@ using System.Linq; using System.Threading.Tasks; using FillInTheTextBot.Models; +using FillInTheTextBot.Services.BackgroundTasks; using FillInTheTextBot.Services.Configuration; using FillInTheTextBot.Services.Extensions; using FillInTheTextBot.Services.Mapping; @@ -17,12 +18,14 @@ public class ConversationService : IConversationService private readonly ConversationConfiguration _configuration; private readonly IDialogflowService _dialogflowService; private readonly IRedisCacheService _cache; + private readonly IBackgroundTaskQueue _backgroundTasks; public ConversationService(ConversationConfiguration configuration, IDialogflowService dialogflowService, - IRedisCacheService cache) + IRedisCacheService cache, IBackgroundTaskQueue backgroundTasks) { _dialogflowService = dialogflowService; _cache = cache; + _backgroundTasks = backgroundTasks; _configuration = configuration; } @@ -289,7 +292,8 @@ private void TrySetSavedText(string sessionId, string scopeKey, Dialog dialog, T { "alternativeText", texts.AlternativeText } }; - _dialogflowService.SetContextAsync(sessionId, scopeKey, "savedText", 5, parameters).Forget(); + _backgroundTasks.Enqueue("savedText context", + () => _dialogflowService.SetContextAsync(sessionId, scopeKey, "savedText", 5, parameters)); } } diff --git a/src/FillInTheTextBot.Services/Extensions/TasksExtensions.cs b/src/FillInTheTextBot.Services/Extensions/TasksExtensions.cs deleted file mode 100644 index b6a9ca43..00000000 --- a/src/FillInTheTextBot.Services/Extensions/TasksExtensions.cs +++ /dev/null @@ -1,37 +0,0 @@ -using Microsoft.Extensions.Logging; -using System; -using System.Threading.Tasks; - -namespace FillInTheTextBot.Services.Extensions -{ - public static class TasksExtensions - { - private static readonly ILogger Log; - - static TasksExtensions() - { - Log = InternalLoggerFactory.CreateLogger(typeof(TaskExtensions).Name); - } - - /// - /// Fire-and-forget - /// Позволяет не дожидаться завершения задачи. - /// В случае ошибки исключение будет логировано. - /// - /// - public static void Forget(this Task task) - { - Task.Factory.StartNew(async () => - { - try - { - await task.ConfigureAwait(false); - } - catch (Exception e) - { - Log?.LogError(e, "Error while executing the task"); - } - }).ConfigureAwait(false); - } - } -} From 7d9a0e0ef7d80fca46f6abbeaf39f3073c75aeaf Mon Sep 17 00:00:00 2001 From: Stepan Grankin Date: Sun, 9 Aug 2026 17:39:54 +0300 Subject: [PATCH 4/6] IBackgroundTaskReader --- .../DI/InternalServicesRegistration.cs | 4 ++++ .../BackgroundTaskProcessorTests.cs | 21 +++++++++--------- .../BackgroundTaskProcessor.cs | 21 +++++------------- .../BackgroundTasks/BackgroundTaskQueue.cs | 5 +---- .../BackgroundTasks/IBackgroundTaskReader.cs | 22 +++++++++++++++++++ 5 files changed, 43 insertions(+), 30 deletions(-) create mode 100644 src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskReader.cs diff --git a/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs b/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs index 03772694..35d3b5e1 100644 --- a/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs +++ b/src/FillInTheTextBot.Api/DI/InternalServicesRegistration.cs @@ -11,8 +11,12 @@ internal static void AddInternalServices(this IServiceCollection services) services.AddTransient(); services.AddScoped(); + // Продюсеры (IBackgroundTaskQueue) и обработчик (IBackgroundTaskReader) должны + // работать с одним и тем же каналом, поэтому оба интерфейса форвардятся на + // единственный экземпляр BackgroundTaskQueue. services.AddSingleton(); services.AddSingleton(provider => provider.GetRequiredService()); + services.AddSingleton(provider => provider.GetRequiredService()); services.AddHostedService(); } } diff --git a/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs index b928c935..4fac3d7f 100644 --- a/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs +++ b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs @@ -6,7 +6,6 @@ using FillInTheTextBot.Services.BackgroundTasks; using Microsoft.Extensions.Logging.Abstractions; using NUnit.Framework; -using NUnit.Framework.Legacy; namespace FillInTheTextBot.Services.Tests.BackgroundTasks { @@ -47,8 +46,8 @@ public async Task Enqueue_Work_Executed() var completed = await Task.WhenAny(executed.Task, Task.Delay(5000)); - ClassicAssert.True(accepted); - ClassicAssert.AreSame(executed.Task, completed, "Работа из очереди должна быть выполнена"); + Assert.That(accepted, Is.True); + Assert.That(completed, Is.SameAs(executed.Task), "Работа из очереди должна быть выполнена"); } [Test] @@ -75,8 +74,8 @@ public async Task Enqueue_ManyWorks_AllExecuted() await WaitForAsync(() => executed.Count == count); - ClassicAssert.AreEqual(count, executed.Count); - CollectionAssert.AreEquivalent(Enumerable.Range(0, count), executed); + Assert.That(executed.Count, Is.EqualTo(count)); + Assert.That(executed, Is.EquivalentTo(Enumerable.Range(0, count))); } [Test] @@ -98,7 +97,7 @@ public async Task Enqueue_FailingWork_QueueKeepsWorking() var completed = await Task.WhenAny(executed.Task, Task.Delay(5000)); - ClassicAssert.AreSame(executed.Task, completed, "Исключение в одной работе не должно останавливать обработчик"); + Assert.That(completed, Is.SameAs(executed.Task), "Исключение в одной работе не должно останавливать обработчик"); } [Test] @@ -132,7 +131,7 @@ public async Task Enqueue_WorkBehindSlowOne_NotBlocked() slowFinish.TrySetResult(true); - ClassicAssert.AreSame(fastExecuted.Task, completed, + Assert.That(completed, Is.SameAs(fastExecuted.Task), "Быстрая работа не должна ждать завершения медленной"); } @@ -155,7 +154,7 @@ public async Task StopAsync_PendingWork_Executed() var completed = await Task.WhenAny(executed.Task, Task.Delay(5000)); - ClassicAssert.AreSame(executed.Task, completed, "Принятые работы должны успеть выполниться при остановке"); + Assert.That(completed, Is.SameAs(executed.Task), "Принятые работы должны успеть выполниться при остановке"); } [Test] @@ -166,14 +165,14 @@ public void Enqueue_QueueIsFull_TaskDropped() { var accepted = _queue.Enqueue($"работа-{i}", () => Task.CompletedTask); - ClassicAssert.True(accepted, $"Работа {i} должна помещаться в очередь"); + Assert.That(accepted, Is.True, $"Работа {i} должна помещаться в очередь"); } var overflow = _queue.Enqueue("лишняя", () => Task.CompletedTask); - ClassicAssert.False(overflow, "Переполнение очереди не должно блокировать вызывающий поток"); + Assert.That(overflow, Is.False, "Переполнение очереди не должно блокировать вызывающий поток"); } [Test] @@ -181,7 +180,7 @@ public void Enqueue_NullWork_NotAccepted() { var accepted = _queue.Enqueue("пустая", null); - ClassicAssert.False(accepted); + Assert.That(accepted, Is.False); } private static async Task WaitForAsync(Func condition, int timeoutMilliseconds = 5000) diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs index ae225062..415fca58 100644 --- a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs +++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs @@ -1,5 +1,4 @@ using System; -using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; @@ -8,7 +7,8 @@ namespace FillInTheTextBot.Services.BackgroundTasks { - public sealed class BackgroundTaskProcessor : BackgroundService + public sealed class BackgroundTaskProcessor(IBackgroundTaskReader queue, ILogger log) + : BackgroundService { /// /// Сколько работ выполняется одновременно. Прежний fire-and-forget запускал их @@ -16,16 +16,7 @@ public sealed class BackgroundTaskProcessor : BackgroundService /// нужно сохранить. Ограничение не даёт при всплеске открыть неограниченное /// число обращений к Redis и Dialogflow. /// - public const int MaxDegreeOfParallelism = 4; - - private readonly BackgroundTaskQueue _queue; - private readonly ILogger _log; - - public BackgroundTaskProcessor(BackgroundTaskQueue queue, ILogger log) - { - _queue = queue; - _log = log; - } + private const int MaxDegreeOfParallelism = 4; protected override Task ExecuteAsync(CancellationToken stoppingToken) { @@ -34,7 +25,7 @@ protected override Task ExecuteAsync(CancellationToken stoppingToken) public override async Task StopAsync(CancellationToken cancellationToken) { - _queue.Complete(); + queue.Complete(); await base.StopAsync(cancellationToken); @@ -64,7 +55,7 @@ private Task DrainAsync() /// private async Task ConsumeAsync() { - await foreach (var task in _queue.ReadAllAsync()) + await foreach (var task in queue.ReadAllAsync()) { await ExecuteTaskAsync(task); } @@ -78,7 +69,7 @@ private async Task ExecuteTaskAsync(BackgroundTask task) } catch (Exception e) { - _log.LogError(e, "Error while executing background task '{TaskName}'", task.Name); + log.LogError(e, "Error while executing background task '{TaskName}'", task.Name); } } } diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs index e00acb1b..0a576cd2 100644 --- a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs +++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs @@ -6,7 +6,7 @@ namespace FillInTheTextBot.Services.BackgroundTasks { - public sealed class BackgroundTaskQueue : IBackgroundTaskQueue + public sealed class BackgroundTaskQueue : IBackgroundTaskQueue, IBackgroundTaskReader { /// /// Ёмкость подобрана с запасом на всплеск запросов: при штатной нагрузке очередь @@ -59,9 +59,6 @@ public IAsyncEnumerable ReadAllAsync() return _channel.Reader.ReadAllAsync(); } - /// - /// Закрывает очередь на запись. Уже принятые работы остаются доступны для чтения. - /// public void Complete() { _channel.Writer.TryComplete(); diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskReader.cs b/src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskReader.cs new file mode 100644 index 00000000..94aba162 --- /dev/null +++ b/src/FillInTheTextBot.Services/BackgroundTasks/IBackgroundTaskReader.cs @@ -0,0 +1,22 @@ +using System.Collections.Generic; + +namespace FillInTheTextBot.Services.BackgroundTasks +{ + /// + /// Потребительская сторона очереди фоновых работ: чтение принятых работ и закрытие + /// очереди на запись при остановке. Отделена от , + /// чтобы обработчик зависел от абстракции, а не от конкретного класса очереди. + /// + public interface IBackgroundTaskReader + { + /// + /// Читает принятые работы, пока очередь не закрыта на запись и не разобрана. + /// + IAsyncEnumerable ReadAllAsync(); + + /// + /// Закрывает очередь на запись. Уже принятые работы остаются доступны для чтения. + /// + void Complete(); + } +} From 4fd84a7a889ecf049468619861591364c9e17670 Mon Sep 17 00:00:00 2001 From: Stepan Grankin Date: Sun, 9 Aug 2026 18:07:10 +0300 Subject: [PATCH 5/6] =?UTF-8?q?=D0=B1=D0=B5=D0=B7=D0=BB=D0=B8=D0=BC=D0=B8?= =?UTF-8?q?=D1=82=D0=BD=D1=8B=D0=B9=20DegreeOfParallelism?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../BackgroundTaskProcessor.cs | 55 ++++++++++--------- 1 file changed, 28 insertions(+), 27 deletions(-) diff --git a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs index 415fca58..18cdea92 100644 --- a/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs +++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs @@ -1,5 +1,5 @@ -using System; -using System.Linq; +using System; +using System.Collections.Concurrent; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.Hosting; @@ -10,14 +10,6 @@ namespace FillInTheTextBot.Services.BackgroundTasks public sealed class BackgroundTaskProcessor(IBackgroundTaskReader queue, ILogger log) : BackgroundService { - /// - /// Сколько работ выполняется одновременно. Прежний fire-and-forget запускал их - /// без ограничений, и медленная работа не задерживала остальные — это свойство - /// нужно сохранить. Ограничение не даёт при всплеске открыть неограниченное - /// число обращений к Redis и Dialogflow. - /// - private const int MaxDegreeOfParallelism = 4; - protected override Task ExecuteAsync(CancellationToken stoppingToken) { return DrainAsync(); @@ -36,29 +28,38 @@ public override async Task StopAsync(CancellationToken cancellationToken) await DrainAsync(); } - private Task DrainAsync() - { - // Несколько независимых потребителей одного канала. Parallel.ForEachAsync здесь - // не подходит: при остановке его задача завершается как отменённая, не дочитав - // очередь, и принятые работы теряются. - var consumers = Enumerable - .Range(0, MaxDegreeOfParallelism) - .Select(_ => ConsumeAsync()); - - return Task.WhenAll(consumers); - } - /// - /// Читает очередь, пока она не закрыта на запись и не разобрана. Отмена по - /// stoppingToken намеренно не используется: при остановке приложения уже принятые - /// работы должны быть доведены до конца, а не потеряны вместе с процессом. + /// Читает очередь и запускает каждую работу, не дожидаясь её завершения: + /// параллелизм не ограничен, как в прежнем fire-and-forget, и медленная или + /// зависшая работа не задерживает остальные. Отмена по stoppingToken намеренно + /// не используется — при остановке приложения уже принятые работы должны быть + /// доведены до конца, поэтому после закрытия очереди метод дожидается запущенных. /// - private async Task ConsumeAsync() + private async Task DrainAsync() { + var running = new ConcurrentDictionary(); + await foreach (var task in queue.ReadAllAsync()) { - await ExecuteTaskAsync(task); + var work = ExecuteTaskAsync(task); + + // Синхронно завершившиеся работы (Task.CompletedTask и т.п.) не трекаем + if (work.IsCompleted) + { + continue; + } + + // Трекаем работы «в полёте», чтобы дождаться их при остановке. Работа + // сама убирает себя из набора по завершении, иначе он рос бы бесконечно. + running[work] = 0; + _ = work.ContinueWith( + finished => running.TryRemove(finished, out _), + CancellationToken.None, + TaskContinuationOptions.ExecuteSynchronously, + TaskScheduler.Default); } + + await Task.WhenAll(running.Keys); } private async Task ExecuteTaskAsync(BackgroundTask task) From 3570e6c6937c4043e7df47451ae900d9d4f674b4 Mon Sep 17 00:00:00 2001 From: Stepan Grankin Date: Sun, 9 Aug 2026 18:21:05 +0300 Subject: [PATCH 6/6] =?UTF-8?q?=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=20=D1=82=D0=B5=D1=81=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../BackgroundTaskProcessorTests.cs | 33 +++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs index 4fac3d7f..b178f68e 100644 --- a/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs +++ b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs @@ -157,6 +157,39 @@ public async Task StopAsync_PendingWork_Executed() Assert.That(completed, Is.SameAs(executed.Task), "Принятые работы должны успеть выполниться при остановке"); } + [Test] + public async Task StopAsync_InFlightWork_WaitsUntilCompleted() + { + var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var finished = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + await _target.StartAsync(CancellationToken.None); + + _queue.Enqueue("в полёте", async () => + { + started.TrySetResult(true); + await release.Task; + finished.TrySetResult(true); + }); + + await started.Task; // работа стартовала и висит на release + + var stop = _target.StopAsync(CancellationToken.None); + + // Пока работа не отпущена, StopAsync не должен завершиться + var early = await Task.WhenAny(stop, Task.Delay(300)); + Assert.That(early, Is.Not.SameAs(stop), "StopAsync не должен завершаться, пока работа в полёте"); + Assert.That(finished.Task.IsCompleted, Is.False); + + release.TrySetResult(true); + + await stop; // после отпускания StopAsync завершается + + Assert.That(finished.Task.IsCompletedSuccessfully, Is.True, + "Работа должна быть доведена до конца до завершения StopAsync"); + } + [Test] public void Enqueue_QueueIsFull_TaskDropped() {