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))
{
// 逐次求值而非启动时一次:本 Worker 无启动延时,宿主配置未必就绪,
// 而它是事件驱动的,空闲时不产生任何查询,逐次求值的代价可以忽略。
if (!AidopJobGate.ShouldRun(nameof(MdpOutboxPushWorker), _logger)) continue;
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");
}
}
}
}