MdpOutboxPushJob.cs 1.6 KB

12345678910111213141516171819202122232425262728293031323334
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. using Furion.Schedule;
  3. using Microsoft.Extensions.DependencyInjection;
  4. using Microsoft.Extensions.Logging;
  5. namespace Admin.NET.Plugin.AiDOP.Job;
  6. /// <summary>
  7. /// 出站 Outbox 兜底扫描:每 60 秒捡漏(进程重启期间入队、触发失败、重试到期)。
  8. /// 主路径为入队后 <see cref="MdpOutboxWakeSignal"/> 事件驱动推送。
  9. /// </summary>
  10. [JobDetail("job_mdp_outbox_push", Description = "MDP Outbox 出站推送(兜底)", GroupName = "default", Concurrent = false)]
  11. [PeriodSeconds(60, TriggerId = "trigger_mdp_outbox_push", Description = "每 60 秒兜底扫描 Outbox", RunOnStart = true)]
  12. public class MdpOutboxPushJob : IJob
  13. {
  14. private readonly IServiceScopeFactory _scopeFactory;
  15. private readonly ILogger _logger;
  16. public MdpOutboxPushJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
  17. {
  18. _scopeFactory = scopeFactory;
  19. _logger = loggerFactory.CreateLogger(nameof(MdpOutboxPushJob));
  20. }
  21. public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
  22. {
  23. using var scope = _scopeFactory.CreateScope();
  24. var dispatcher = scope.ServiceProvider.GetRequiredService<MdpTargetPushDispatcher>();
  25. var (success, failed, skipped) = await dispatcher.PushPendingAsync(
  26. MdpTargetPushDispatcher.DefaultTake, stoppingToken);
  27. if (success + failed + skipped > 0)
  28. _logger.LogInformation("[MdpOutboxPushJob] success={Success} failed={Failed} retrying={Skipped}", success, failed, skipped);
  29. }
  30. }