using Admin.NET.Core;
using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
using Furion.Schedule;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
namespace Admin.NET.Plugin.AiDOP.Job;
///
/// WP10 S4b / D12:Outbox 死信与积压告警(LogWarning + 站内 SysNotice)。
/// 每 5 分钟扫描;同条件告警间隔不少于 30 分钟,避免刷屏。无钉钉。
///
[JobDetail("job_mdp_outbox_dl_alert", Description = "MDP Outbox 死信/积压告警", GroupName = "default", Concurrent = false)]
[PeriodSeconds(300, TriggerId = "trigger_mdp_outbox_dl_alert", Description = "每 5 分钟扫描 Outbox 告警", RunOnStart = false)]
public class MdpOutboxDeadLetterAlertJob : IJob
{
private const long NoticeReceiverUserId = 1300000000101L;
private const string NoticeReceiverUserName = "超级管理员";
private static readonly TimeSpan AlertCooldown = TimeSpan.FromMinutes(30);
private static readonly TimeSpan StalePendingThreshold = TimeSpan.FromMinutes(15);
private static readonly object Gate = new();
private static DateTime? _lastAlertUtc;
private static DateTime? _lastScanUtc;
private static long _lastSeenDeadMaxId;
private readonly IServiceScopeFactory _scopeFactory;
private readonly ILogger _logger;
public MdpOutboxDeadLetterAlertJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
{
_scopeFactory = scopeFactory;
_logger = loggerFactory.CreateLogger(nameof(MdpOutboxDeadLetterAlertJob));
}
public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
{
using var scope = _scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
var now = DateTime.Now;
var scanSince = _lastScanUtc?.ToLocalTime() ?? now.AddMinutes(-5);
var newDeadCount = await db.Queryable()
.CountAsync(x => x.Status == 2 && x.UpdateTime >= scanSince, stoppingToken);
var maxDeadId = await db.Queryable()
.Where(x => x.Status == 2)
.MaxAsync(x => (long?)x.Id, stoppingToken) ?? 0L;
var newDeadById = maxDeadId > _lastSeenDeadMaxId
? await db.Queryable()
.CountAsync(x => x.Status == 2 && x.Id > _lastSeenDeadMaxId, stoppingToken)
: 0;
var hasNewDead = newDeadCount > 0 || newDeadById > 0;
double oldestPendingMinutes = 0;
var pending = await db.Queryable().CountAsync(x => x.Status == 0, stoppingToken);
if (pending > 0)
{
var oldest = await db.Queryable()
.Where(x => x.Status == 0)
.OrderBy(x => x.CreateTime)
.Select(x => x.CreateTime)
.FirstAsync(stoppingToken);
oldestPendingMinutes = (now - oldest).TotalMinutes;
}
var stalePending = oldestPendingMinutes > StalePendingThreshold.TotalMinutes;
_lastScanUtc = DateTime.UtcNow;
if (maxDeadId > _lastSeenDeadMaxId)
_lastSeenDeadMaxId = maxDeadId;
if (!hasNewDead && !stalePending)
return;
lock (Gate)
{
if (_lastAlertUtc.HasValue && DateTime.UtcNow - _lastAlertUtc.Value < AlertCooldown)
return;
_lastAlertUtc = DateTime.UtcNow;
}
var deadTotal = await db.Queryable().CountAsync(x => x.Status == 2, stoppingToken);
var msg =
$"Outbox 告警:新增死信≈{Math.Max(newDeadCount, newDeadById)}(累计死信={deadTotal})," +
$"待推={pending},最老待推≈{oldestPendingMinutes:F1} 分钟。" +
$"请到「出站回写队列」查看 /aidop/data-platform/outbox。";
_logger.LogWarning("[MdpOutboxDeadLetterAlertJob] {Message}", msg);
try
{
// 后台作业无登录态,直接落 SysNotice + 一条 SysNoticeUser(超管),不走 SysNoticeService.InitNoticeInfo
var notice = new SysNotice
{
Title = "Outbox死信/积压告警",
Content = msg,
Type = NoticeTypeEnum.NOTICE,
PublicUserId = NoticeReceiverUserId,
PublicUserName = NoticeReceiverUserName,
PublicTime = now,
Status = NoticeStatusEnum.PUBLIC,
CreateTime = now,
CreateUserId = NoticeReceiverUserId,
CreateUserName = NoticeReceiverUserName,
};
var noticeId = await db.Insertable(notice).ExecuteReturnSnowflakeIdAsync(stoppingToken);
await db.Insertable(new SysNoticeUser
{
NoticeId = noticeId,
UserId = NoticeReceiverUserId,
ReadStatus = NoticeUserStatusEnum.UNREAD,
}).ExecuteCommandAsync(stoppingToken);
}
catch (Exception ex)
{
_logger.LogError(ex, "[MdpOutboxDeadLetterAlertJob] SysNotice 写入失败");
}
}
}