MdpOutboxDeadLetterAlertJob.cs 5.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125
  1. using Admin.NET.Core;
  2. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  3. using Furion.Schedule;
  4. using Microsoft.Extensions.DependencyInjection;
  5. using Microsoft.Extensions.Logging;
  6. namespace Admin.NET.Plugin.AiDOP.Job;
  7. /// <summary>
  8. /// WP10 S4b / D12:Outbox 死信与积压告警(LogWarning + 站内 SysNotice)。
  9. /// 每 5 分钟扫描;同条件告警间隔不少于 30 分钟,避免刷屏。无钉钉。
  10. /// </summary>
  11. [JobDetail("job_mdp_outbox_dl_alert", Description = "MDP Outbox 死信/积压告警", GroupName = "default", Concurrent = false)]
  12. [PeriodSeconds(300, TriggerId = "trigger_mdp_outbox_dl_alert", Description = "每 5 分钟扫描 Outbox 告警", RunOnStart = false)]
  13. public class MdpOutboxDeadLetterAlertJob : IJob
  14. {
  15. private const long NoticeReceiverUserId = 1300000000101L;
  16. private const string NoticeReceiverUserName = "超级管理员";
  17. private static readonly TimeSpan AlertCooldown = TimeSpan.FromMinutes(30);
  18. private static readonly TimeSpan StalePendingThreshold = TimeSpan.FromMinutes(15);
  19. private static readonly object Gate = new();
  20. private static DateTime? _lastAlertUtc;
  21. private static DateTime? _lastScanUtc;
  22. private static long _lastSeenDeadMaxId;
  23. private readonly IServiceScopeFactory _scopeFactory;
  24. private readonly ILogger _logger;
  25. public MdpOutboxDeadLetterAlertJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
  26. {
  27. _scopeFactory = scopeFactory;
  28. _logger = loggerFactory.CreateLogger(nameof(MdpOutboxDeadLetterAlertJob));
  29. }
  30. public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
  31. {
  32. if (!AidopJobGate.ShouldRun(nameof(MdpOutboxDeadLetterAlertJob), _logger)) return;
  33. using var scope = _scopeFactory.CreateScope();
  34. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  35. var now = DateTime.Now;
  36. var scanSince = _lastScanUtc?.ToLocalTime() ?? now.AddMinutes(-5);
  37. var newDeadCount = await db.Queryable<MdpOutbox>()
  38. .CountAsync(x => x.Status == 2 && x.UpdateTime >= scanSince, stoppingToken);
  39. var maxDeadId = await db.Queryable<MdpOutbox>()
  40. .Where(x => x.Status == 2)
  41. .MaxAsync(x => (long?)x.Id, stoppingToken) ?? 0L;
  42. // SqlSugar 表达式树不能取私有静态字段的值,先落到局部变量再进 lambda
  43. var lastSeenDeadMaxId = _lastSeenDeadMaxId;
  44. var newDeadById = maxDeadId > lastSeenDeadMaxId
  45. ? await db.Queryable<MdpOutbox>()
  46. .CountAsync(x => x.Status == 2 && x.Id > lastSeenDeadMaxId, stoppingToken)
  47. : 0;
  48. var hasNewDead = newDeadCount > 0 || newDeadById > 0;
  49. double oldestPendingMinutes = 0;
  50. var pending = await db.Queryable<MdpOutbox>().CountAsync(x => x.Status == 0, stoppingToken);
  51. if (pending > 0)
  52. {
  53. var oldest = await db.Queryable<MdpOutbox>()
  54. .Where(x => x.Status == 0)
  55. .OrderBy(x => x.CreateTime)
  56. .Select(x => x.CreateTime)
  57. .FirstAsync(stoppingToken);
  58. oldestPendingMinutes = (now - oldest).TotalMinutes;
  59. }
  60. var stalePending = oldestPendingMinutes > StalePendingThreshold.TotalMinutes;
  61. _lastScanUtc = DateTime.UtcNow;
  62. if (maxDeadId > lastSeenDeadMaxId)
  63. _lastSeenDeadMaxId = maxDeadId;
  64. if (!hasNewDead && !stalePending)
  65. return;
  66. lock (Gate)
  67. {
  68. if (_lastAlertUtc.HasValue && DateTime.UtcNow - _lastAlertUtc.Value < AlertCooldown)
  69. return;
  70. _lastAlertUtc = DateTime.UtcNow;
  71. }
  72. var deadTotal = await db.Queryable<MdpOutbox>().CountAsync(x => x.Status == 2, stoppingToken);
  73. var msg =
  74. $"Outbox 告警:新增死信≈{Math.Max(newDeadCount, newDeadById)}(累计死信={deadTotal})," +
  75. $"待推={pending},最老待推≈{oldestPendingMinutes:F1} 分钟。" +
  76. $"请到「出站回写队列」查看 /aidop/data-platform/outbox。";
  77. _logger.LogWarning("[MdpOutboxDeadLetterAlertJob] {Message}", msg);
  78. try
  79. {
  80. // 后台作业无登录态,直接落 SysNotice + 一条 SysNoticeUser(超管),不走 SysNoticeService.InitNoticeInfo
  81. var notice = new SysNotice
  82. {
  83. Title = "Outbox死信/积压告警",
  84. Content = msg,
  85. Type = NoticeTypeEnum.NOTICE,
  86. PublicUserId = NoticeReceiverUserId,
  87. PublicUserName = NoticeReceiverUserName,
  88. PublicTime = now,
  89. Status = NoticeStatusEnum.PUBLIC,
  90. CreateTime = now,
  91. CreateUserId = NoticeReceiverUserId,
  92. CreateUserName = NoticeReceiverUserName,
  93. };
  94. var noticeId = await db.Insertable(notice).ExecuteReturnSnowflakeIdAsync(stoppingToken);
  95. await db.Insertable(new SysNoticeUser
  96. {
  97. NoticeId = noticeId,
  98. UserId = NoticeReceiverUserId,
  99. ReadStatus = NoticeUserStatusEnum.UNREAD,
  100. }).ExecuteCommandAsync(stoppingToken);
  101. }
  102. catch (Exception ex)
  103. {
  104. _logger.LogError(ex, "[MdpOutboxDeadLetterAlertJob] SysNotice 写入失败");
  105. }
  106. }
  107. }