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");
}
}
}
}