/
vshmidt
/
masstransit
Обзор
Документация
Войти
/
vshmidt
/
masstransit
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
develop
src/Transports/MassTransit.EventHubIntegration/EventHubIntegration/ProcessorLockContext.cs
91 строка
3 KB
Chris Patterson
Fixed #6049 - avoid failure when reusing the client on a subsequent receive transport restart
25 июл 2025, 18:49
25 июл 2025, 18:49
9f46a6c
Код
Авторство
О чём код?
namespace MassTransit.EventHubIntegration { using System; using System.Threading; using System.Threading.Tasks; using Azure.Messaging.EventHubs; using Azure.Messaging.EventHubs.Processor; using Checkpoints; using Util; public class ProcessorLockContext : IProcessorLockContext, ProcessorClientBuilderContext { readonly ProcessorContext _context; readonly SingleThreadedDictionary<string, PartitionCheckpointData> _data; readonly PendingConfirmationCollection _pending; readonly ReceiveSettings _receiveSettings; public ProcessorLockContext(ProcessorContext context, ReceiveSettings receiveSettings, CancellationToken cancellationToken) { _context = context; _receiveSettings = receiveSettings; _pending = new PendingConfirmationCollection(cancellationToken); _data = new SingleThreadedDictionary<string, PartitionCheckpointData>(StringComparer.Ordinal); Client = context.GetClient(this); } public EventProcessorClient Client { get; } public ValueTask DisposeAsync() { _context.ReleaseClient(this); _pending.Dispose(); return default; } public Task Pending(ProcessEventArgs eventArgs) { LogContext.SetCurrentIfNull(_context.LogContext); return _data.TryGetValue(eventArgs.Partition.PartitionId, out var data) ? data.Pending(eventArgs) : Task.CompletedTask; } public Task Faulted(ProcessEventArgs eventArgs, Exception exception) { LogContext.SetCurrentIfNull(_context.LogContext); _pending.Faulted(eventArgs, exception); return Task.CompletedTask; } public Task Complete(ProcessEventArgs eventArgs) { LogContext.SetCurrentIfNull(_context.LogContext); _pending.Complete(eventArgs); return Task.CompletedTask; } public void Canceled(ProcessEventArgs eventArgs, CancellationToken cancellationToken) { LogContext.SetCurrentIfNull(_context.LogContext); _pending.Canceled(eventArgs, cancellationToken); } public Task OnPartitionInitializing(PartitionInitializingEventArgs eventArgs) { LogContext.SetCurrentIfNull(_context.LogContext); if (_data.TryAdd(eventArgs.PartitionId, _ => new PartitionCheckpointData(_receiveSettings, _pending))) LogContext.Info?.Log("Partition: {PartitionId} was initialized", eventArgs.PartitionId); return Task.CompletedTask; } public Task OnPartitionClosing(PartitionClosingEventArgs eventArgs) { LogContext.SetCurrentIfNull(_context.LogContext); return _data.TryRemove(eventArgs.PartitionId, out var data) ? data.Close(eventArgs) : Task.CompletedTask; } } }