소스 검색

add failed message processor.

master
Savorboard 7 년 전
부모
커밋
8e41e0d552
1개의 변경된 파일3개의 추가작업 그리고 6개의 파일을 삭제
  1. +3
    -6
      src/DotNetCore.CAP/Processor/IProcessingServer.Cap.cs

+ 3
- 6
src/DotNetCore.CAP/Processor/IProcessingServer.Cap.cs 파일 보기

@@ -47,7 +47,7 @@ namespace DotNetCore.CAP.Processor
_context = new ProcessingContext(_provider, _cts.Token);

var processorTasks = _processors
.Select(p => InfiniteRetry(p))
.Select(InfiniteRetry)
.Select(p => p.ProcessAsync(_context));
_compositeTask = Task.WhenAll(processorTasks);
}
@@ -84,10 +84,7 @@ namespace DotNetCore.CAP.Processor

private bool AllProcessorsWaiting()
{
foreach (var processor in _messageDispatchers)
if (!processor.Waiting)
return false;
return true;
return _messageDispatchers.All(processor => processor.Waiting);
}

private IProcessor InfiniteRetry(IProcessor inner)
@@ -107,7 +104,7 @@ namespace DotNetCore.CAP.Processor

returnedProcessors.Add(_provider.GetRequiredService<PublishQueuer>());
returnedProcessors.Add(_provider.GetRequiredService<SubscribeQueuer>());
//returnedProcessors.Add(_provider.GetRequiredService<FailedJobProcessor>());
returnedProcessors.Add(_provider.GetRequiredService<FailedProcessor>());

returnedProcessors.Add(_provider.GetRequiredService<IAdditionalProcessor>());



불러오는 중...
취소
저장