/
vshmidt
/
masstransit
Обзор
Документация
Войти
/
vshmidt
/
masstransit
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
develop
src/MassTransit/Testing/Implementations/AsyncInactivityObserver.cs
82 строки
2 KB
Chris Patterson
Inspired by #3619 - Reworked the way inactivity timers work to more consistently handle timer resets as messages are received.
19 авг 2022, 17:57
19 авг 2022, 17:57
9d6e0cb
Код
Авторство
О чём код?
namespace MassTransit.Testing.Implementations { using System; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; using Internals; using Util; public class AsyncInactivityObserver : IInactivityObserver { readonly Lazy<Task> _inactivityTask; readonly TaskCompletionSource<bool> _inactivityTaskSource; readonly CancellationTokenSource _inactivityTokenSource; readonly HashSet<IInactivityObservationSource> _sources; public AsyncInactivityObserver(TimeSpan timeout, CancellationToken cancellationToken) { _inactivityTaskSource = TaskUtil.GetTask(); _inactivityTask = new Lazy<Task>(() => TimeoutTask(timeout, cancellationToken)); _sources = new HashSet<IInactivityObservationSource>(); _inactivityTokenSource = new CancellationTokenSource(); } public Task InactivityTask => _inactivityTask.Value; public CancellationToken InactivityToken => _inactivityTokenSource.Token; public void Connected(IInactivityObservationSource source) { _sources.Add(source); } public Task NoActivity() { return CheckSourceActivity(); } public void ForceInactive() { _inactivityTaskSource.TrySetResult(true); _inactivityTokenSource.Cancel(); } Task<bool> CheckSourceActivity() { if (_sources.All(x => x.IsInactive)) { _inactivityTaskSource.TrySetResult(true); _inactivityTokenSource.Cancel(); return TaskUtil.True; } return TaskUtil.False; } async Task TimeoutTask(TimeSpan timeout, CancellationToken cancellationToken) { try { var inActive = false; do { await Task.Delay(timeout, cancellationToken).ConfigureAwait(false); inActive = await CheckSourceActivity().ConfigureAwait(false); } while (!inActive); await _inactivityTaskSource.Task.OrCanceled(cancellationToken).ConfigureAwait(false); } catch (Exception) { } } } }