using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors; /// 事件驱动出站:收到 后立即推一批。 public sealed class MdpOutboxPushWorker : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; private readonly MdpOutboxWakeSignal _wake; private readonly ILogger _logger; public MdpOutboxPushWorker( IServiceScopeFactory scopeFactory, MdpOutboxWakeSignal wake, ILoggerFactory loggerFactory) { _scopeFactory = scopeFactory; _wake = wake; _logger = loggerFactory.CreateLogger(nameof(MdpOutboxPushWorker)); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await foreach (var _ in _wake.Reader.ReadAllAsync(stoppingToken)) { try { using var scope = _scopeFactory.CreateScope(); var dispatcher = scope.ServiceProvider.GetRequiredService(); var (success, failed, skipped) = await dispatcher.PushPendingAsync( MdpTargetPushDispatcher.DefaultTake, stoppingToken); if (success + failed + skipped > 0) _logger.LogInformation( "[MdpOutboxPushWorker] success={Success} failed={Failed} retrying={Skipped}", success, failed, skipped); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { break; } catch (Exception ex) { _logger.LogWarning(ex, "[MdpOutboxPushWorker] push failed"); } } } }