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..35d3b5e1 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,14 @@ 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.Messengers.Marusia.Tests/MarusiaServiceTests.cs b/src/FillInTheTextBot.Messengers.Marusia.Tests/MarusiaServiceTests.cs
index a6f2a234..c0f9ccef 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;
@@ -22,6 +23,7 @@ public class MarusiaServiceTests
{
private Mock _conversationService;
private Mock _cache;
+ private Mock _backgroundTasks;
private MarusiaService _target;
@@ -33,7 +35,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 b3d93d80..cc89ce25 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;
@@ -20,6 +21,7 @@ public class SberServiceTests
{
private Mock _conversationService;
private Mock _cache;
+ private Mock _backgroundTasks;
private SberService _target;
@@ -31,7 +33,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/BackgroundTasks/BackgroundTaskProcessorTests.cs b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs
new file mode 100644
index 00000000..b178f68e
--- /dev/null
+++ b/src/FillInTheTextBot.Services.Tests/BackgroundTasks/BackgroundTaskProcessorTests.cs
@@ -0,0 +1,230 @@
+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;
+
+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));
+
+
+ Assert.That(accepted, Is.True);
+ Assert.That(completed, Is.SameAs(executed.Task), "Работа из очереди должна быть выполнена");
+ }
+
+ [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);
+
+
+ Assert.That(executed.Count, Is.EqualTo(count));
+ Assert.That(executed, Is.EquivalentTo(Enumerable.Range(0, count)));
+ }
+
+ [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));
+
+
+ Assert.That(completed, Is.SameAs(executed.Task), "Исключение в одной работе не должно останавливать обработчик");
+ }
+
+ [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);
+
+ Assert.That(completed, Is.SameAs(fastExecuted.Task),
+ "Быстрая работа не должна ждать завершения медленной");
+ }
+
+ [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));
+
+ 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()
+ {
+ // Обработчик не запущен, поэтому очередь только наполняется
+ foreach (var i in Enumerable.Range(0, BackgroundTaskQueue.Capacity))
+ {
+ var accepted = _queue.Enqueue($"работа-{i}", () => Task.CompletedTask);
+
+ Assert.That(accepted, Is.True, $"Работа {i} должна помещаться в очередь");
+ }
+
+
+ var overflow = _queue.Enqueue("лишняя", () => Task.CompletedTask);
+
+
+ Assert.That(overflow, Is.False, "Переполнение очереди не должно блокировать вызывающий поток");
+ }
+
+ [Test]
+ public void Enqueue_NullWork_NotAccepted()
+ {
+ var accepted = _queue.Enqueue("пустая", null);
+
+ Assert.That(accepted, Is.False);
+ }
+
+ 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.Tests/ConversationServiceTests.cs b/src/FillInTheTextBot.Services.Tests/ConversationServiceTests.cs
index f5dabc8b..a4aca0e8 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;
@@ -19,6 +21,7 @@ public class ConversationServiceTests
{
private Mock _dialogflowService;
private Mock _cache;
+ private Mock _backgroundTasks;
private ConversationConfiguration _configuration;
@@ -32,9 +35,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/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..18cdea92
--- /dev/null
+++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskProcessor.cs
@@ -0,0 +1,77 @@
+using System;
+using System.Collections.Concurrent;
+using System.Threading;
+using System.Threading.Tasks;
+using Microsoft.Extensions.Hosting;
+using Microsoft.Extensions.Logging;
+
+namespace FillInTheTextBot.Services.BackgroundTasks
+{
+ public sealed class BackgroundTaskProcessor(IBackgroundTaskReader queue, ILogger log)
+ : BackgroundService
+ {
+ protected override Task ExecuteAsync(CancellationToken stoppingToken)
+ {
+ return DrainAsync();
+ }
+
+ public override async Task StopAsync(CancellationToken cancellationToken)
+ {
+ queue.Complete();
+
+ await base.StopAsync(cancellationToken);
+
+ // BackgroundService запускает ExecuteAsync отложенно, и при остановке вскоре
+ // после старта задача отменяется до входа в тело метода — тогда очередь никто
+ // не разобрал. Поэтому остаток добирается здесь, независимо от того,
+ // успел ли стартовать основной цикл.
+ await DrainAsync();
+ }
+
+ ///
+ /// Читает очередь и запускает каждую работу, не дожидаясь её завершения:
+ /// параллелизм не ограничен, как в прежнем fire-and-forget, и медленная или
+ /// зависшая работа не задерживает остальные. Отмена по stoppingToken намеренно
+ /// не используется — при остановке приложения уже принятые работы должны быть
+ /// доведены до конца, поэтому после закрытия очереди метод дожидается запущенных.
+ ///
+ private async Task DrainAsync()
+ {
+ var running = new ConcurrentDictionary();
+
+ await foreach (var task in queue.ReadAllAsync())
+ {
+ 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)
+ {
+ 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..0a576cd2
--- /dev/null
+++ b/src/FillInTheTextBot.Services/BackgroundTasks/BackgroundTaskQueue.cs
@@ -0,0 +1,67 @@
+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, IBackgroundTaskReader
+ {
+ ///
+ /// Ёмкость подобрана с запасом на всплеск запросов: при штатной нагрузке очередь
+ /// разбирается быстрее, чем наполняется, а при недоступности внешнего сервиса
+ /// ограничение не даёт очереди съесть память.
+ ///
+ 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/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();
+ }
+}
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);
- }
- }
-}
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 @@
+