Files
IOT_API/Service/Implement/Inspection/AlertNotifyService.cs
T

592 lines
23 KiB
C#

using Model;
using Model.Dto.Inspection;
using Model.Entity.Inspection;
using Model.Mapper;
using ORM;
using Service.Interface;
using SqlSugar;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
namespace Service.Implement
{
/// <summary>
/// 告警通知配置服务实现(渠道/通知规则/推送日志/值班排班/跨网发件箱)
/// </summary>
public class AlertNotifyService : IAlertNotifyService
{
private readonly IEnumerable<IAlertNotifier> _notifiers;
public AlertNotifyService(IEnumerable<IAlertNotifier> notifiers)
{
_notifiers = notifiers;
}
#region 通知渠道
/// <summary>
/// 渠道列表(Secret 脱敏,只返回 HasSecret 标记)
/// </summary>
public async Task<Result<List<AlertNotifyChannelDto>>> GetChannelsAsync()
{
try
{
var list = await SqlSugarContext.DbContext.Queryable<AlertNotifyChannelEntity>()
.Where(x => x.IsDel == 0)
.OrderBy(x => x.CreateTime, OrderByType.Desc)
.ToListAsync();
return Result<List<AlertNotifyChannelDto>>.Success(list.ToDtoList());
}
catch (Exception ex)
{
return Result<List<AlertNotifyChannelDto>>.Error("查询通知渠道失败", ex);
}
}
/// <summary>
/// 新增渠道
/// </summary>
public async Task<Result<bool>> AddChannelAsync(AlertNotifyChannelDto dto)
{
if (string.IsNullOrWhiteSpace(dto?.Name) || string.IsNullOrWhiteSpace(dto?.WebhookUrl))
{
return Result<bool>.Error("渠道名称与 Webhook 地址不能为空");
}
try
{
var entity = dto.ToEntity();
entity.Id = 0; // 防止前端误传 Id,新增一律由雪花生成
entity.CreateTime = DateTime.Now;
await SqlSugarContext.DbContext.Insertable(entity).ExecuteCommandAsync();
return Result<bool>.Success(true);
}
catch (Exception ex)
{
return Result<bool>.Error("新增通知渠道失败", ex);
}
}
/// <summary>
/// 修改渠道(Secret 传空表示保持原值不覆盖——出口已脱敏,前端编辑时不回传明文)
/// </summary>
public async Task<Result<bool>> UpdateChannelAsync(AlertNotifyChannelDto dto)
{
var entity = dto?.ToEntity();
if (entity == null || entity.Id <= 0)
{
return Result<bool>.Error("渠道 Id 无效");
}
try
{
var updater = SqlSugarContext.DbContext.Updateable<AlertNotifyChannelEntity>()
.SetColumns(x => new AlertNotifyChannelEntity
{
ChannelType = entity.ChannelType,
Name = entity.Name,
WebhookUrl = entity.WebhookUrl,
PushMode = entity.PushMode,
IsEnabled = entity.IsEnabled,
Remark = entity.Remark
});
// Secret 显式传值时才覆盖(空表示不修改)
if (!string.IsNullOrWhiteSpace(entity.Secret))
{
updater = updater.SetColumns(x => x.Secret == entity.Secret);
}
var rows = await updater.Where(x => x.Id == entity.Id && x.IsDel == 0).ExecuteCommandAsync();
return Result<bool>.Success(rows > 0);
}
catch (Exception ex)
{
return Result<bool>.Error("修改通知渠道失败", ex);
}
}
/// <summary>
/// 删除渠道(软删除)
/// </summary>
public async Task<Result<bool>> DeleteChannelAsync(long id)
{
try
{
var rows = await SqlSugarContext.DbContext.Updateable<AlertNotifyChannelEntity>()
.SetColumns(x => x.IsDel == 1)
.Where(x => x.Id == id)
.ExecuteCommandAsync();
return Result<bool>.Success(rows > 0);
}
catch (Exception ex)
{
return Result<bool>.Error("删除通知渠道失败", ex);
}
}
/// <summary>
/// 渠道测试推送(发送测试消息验证 Webhook 配置,结果写入推送日志)
/// </summary>
public async Task<Result<bool>> TestSendAsync(long channelId)
{
try
{
var channel = await SqlSugarContext.DbContext.Queryable<AlertNotifyChannelEntity>()
.Where(x => x.Id == channelId && x.IsDel == 0)
.FirstAsync();
if (channel == null)
{
return Result<bool>.Error("渠道不存在或已被删除");
}
var notifier = _notifiers.FirstOrDefault(n => n.ChannelType == channel.ChannelType);
if (notifier == null)
{
return Result<bool>.Error($"暂不支持的渠道类型: {channel.ChannelType}");
}
var message = new AlertNotifyMessage
{
Title = "【测试】设备告警通知",
Content = $"这是一条来自 IOT 设备管理平台的测试消息。\n渠道: {channel.Name}\n时间: {DateTime.Now:yyyy-MM-dd HH:mm:ss}"
};
NotifyResult sendResult;
if (channel.PushMode == NotifyPushModeEnum.Outbox)
{
// Outbox 模式:测试消息也走发件箱,由 AlertBridge 实际推送
await InsertOutboxAsync(0, channel, message, new List<string>());
sendResult = NotifyResult.Ok();
}
else
{
sendResult = await notifier.SendAsync(message, channel.WebhookUrl, channel.Secret);
}
await InsertLogAsync(0, channel, message.Title + " " + message.Content, sendResult, 0);
return sendResult.Success
? Result<bool>.Success(true)
: Result<bool>.Error($"测试推送失败: {sendResult.ErrorMsg}");
}
catch (Exception ex)
{
return Result<bool>.Error("测试推送失败", ex);
}
}
#endregion
#region 通知规则
/// <summary>
/// 规则列表
/// </summary>
public async Task<Result<List<AlertNotifyRuleDto>>> GetRulesAsync()
{
try
{
var list = await SqlSugarContext.DbContext.Queryable<AlertNotifyRuleEntity>()
.Where(x => x.IsDel == 0)
.OrderBy(x => x.AlertLevel, OrderByType.Desc)
.OrderBy(x => x.CreateTime, OrderByType.Desc)
.ToListAsync();
return Result<List<AlertNotifyRuleDto>>.Success(list.ToDtoList());
}
catch (Exception ex)
{
return Result<List<AlertNotifyRuleDto>>.Error("查询通知规则失败", ex);
}
}
/// <summary>
/// 新增规则
/// </summary>
public async Task<Result<bool>> AddRuleAsync(AlertNotifyRuleDto dto)
{
if (string.IsNullOrWhiteSpace(dto?.RuleName))
{
return Result<bool>.Error("规则名称不能为空");
}
if (dto.ChannelId == null || !long.TryParse(dto.ChannelId, out var channelId) || channelId <= 0)
{
return Result<bool>.Error("请选择有效的通知渠道");
}
try
{
var entity = dto.ToEntity();
entity.Id = 0;
entity.CreateTime = DateTime.Now;
await SqlSugarContext.DbContext.Insertable(entity).ExecuteCommandAsync();
return Result<bool>.Success(true);
}
catch (Exception ex)
{
return Result<bool>.Error("新增通知规则失败", ex);
}
}
/// <summary>
/// 修改规则
/// </summary>
public async Task<Result<bool>> UpdateRuleAsync(AlertNotifyRuleDto dto)
{
var entity = dto?.ToEntity();
if (entity == null || entity.Id <= 0)
{
return Result<bool>.Error("规则 Id 无效");
}
try
{
var rows = await SqlSugarContext.DbContext.Updateable(entity)
.IgnoreColumns(x => new { x.CreateTime, x.IsDel })
.ExecuteCommandAsync();
return Result<bool>.Success(rows > 0);
}
catch (Exception ex)
{
return Result<bool>.Error("修改通知规则失败", ex);
}
}
/// <summary>
/// 删除规则(软删除)
/// </summary>
public async Task<Result<bool>> DeleteRuleAsync(long id)
{
try
{
var rows = await SqlSugarContext.DbContext.Updateable<AlertNotifyRuleEntity>()
.SetColumns(x => x.IsDel == 1)
.Where(x => x.Id == id)
.ExecuteCommandAsync();
return Result<bool>.Success(rows > 0);
}
catch (Exception ex)
{
return Result<bool>.Error("删除通知规则失败", ex);
}
}
/// <summary>
/// 启用/停用规则
/// </summary>
public async Task<Result<bool>> SetRuleEnabledAsync(long id, bool enabled)
{
try
{
var rows = await SqlSugarContext.DbContext.Updateable<AlertNotifyRuleEntity>()
.SetColumns(x => x.IsEnabled == enabled)
.Where(x => x.Id == id)
.ExecuteCommandAsync();
return Result<bool>.Success(rows > 0);
}
catch (Exception ex)
{
return Result<bool>.Error("更新规则启用状态失败", ex);
}
}
#endregion
#region 推送日志
/// <summary>
/// 推送日志分页(可按成功状态/渠道筛选)
/// </summary>
public async Task<Result<List<AlertNotifyLogDto>>> GetLogsPagedAsync(int pageIndex, int pageSize, RefAsync<int> total,
bool? success, long channelId = 0)
{
try
{
var list = await SqlSugarContext.DbContext.Queryable<AlertNotifyLogEntity>()
.Where(x => x.IsDel == 0)
.WhereIF(success.HasValue, x => x.Success == success!.Value)
.WhereIF(channelId > 0, x => x.ChannelId == channelId)
.OrderBy(x => x.NotifyTime, OrderByType.Desc)
.ToPageListAsync(pageIndex, pageSize, total);
return Result<List<AlertNotifyLogDto>>.Success(list.ToDtoList());
}
catch (Exception ex)
{
return Result<List<AlertNotifyLogDto>>.Error("查询推送日志失败", ex);
}
}
/// <summary>
/// 写入推送日志(内部公共方法,Worker 也复用)
/// </summary>
internal async Task InsertLogAsync(long alertId, AlertNotifyChannelEntity channel, string content, NotifyResult result, int retryCount)
{
try
{
await SqlSugarContext.DbContext.Insertable(new AlertNotifyLogEntity
{
AlertId = alertId,
ChannelType = channel.ChannelType,
ChannelId = channel.Id,
Target = MaskUrl(channel.WebhookUrl),
Content = content.Length > 1000 ? content.Substring(0, 1000) : content,
Success = result.Success,
ErrorMsg = result.ErrorMsg,
RetryCount = retryCount,
NotifyTime = DateTime.Now,
CreateTime = DateTime.Now
}).ExecuteCommandAsync();
}
catch
{
// 日志写入失败不影响主流程
}
}
/// <summary>
/// Webhook 地址脱敏(保留主机名与路径结构,隐藏末段令牌与查询参数)
/// 飞书令牌在路径末段(/hook/{token})、钉钉在 access_token、企微在 key,均需遮蔽,避免日志泄露密钥
/// </summary>
internal static string MaskUrl(string url)
{
if (string.IsNullOrWhiteSpace(url)) return url;
try
{
var uri = new Uri(url);
var segments = uri.AbsolutePath.Split('/', StringSplitOptions.RemoveEmptyEntries);
// 保留除末段外的路径,末段(如飞书 hook 令牌)用 *** 代替;查询串一律丢弃
string maskedPath = segments.Length > 1
? "/" + string.Join("/", segments.Take(segments.Length - 1)) + "/***"
: "/***";
return $"{uri.Scheme}://{uri.Host}{maskedPath}";
}
catch
{
return url.Length > 60 ? url.Substring(0, 60) + "..." : url;
}
}
#endregion
#region 值班排班
/// <summary>
/// 值班排班列表(按日期范围,缺省查本月)
/// </summary>
public async Task<Result<List<DutyRosterDto>>> GetDutiesAsync(DateTime? startDate, DateTime? endDate)
{
try
{
var start = startDate ?? new DateTime(DateTime.Today.Year, DateTime.Today.Month, 1);
var end = endDate ?? start.AddMonths(1).AddDays(-1);
var list = await SqlSugarContext.DbContext.Queryable<DutyRosterEntity>()
.Where(x => x.IsDel == 0 && x.DutyDate >= start.Date && x.DutyDate <= end.Date)
.OrderBy(x => x.DutyDate, OrderByType.Asc)
.ToListAsync();
return Result<List<DutyRosterDto>>.Success(list.ToDtoList());
}
catch (Exception ex)
{
return Result<List<DutyRosterDto>>.Error("查询值班排班失败", ex);
}
}
/// <summary>
/// 新增排班
/// </summary>
public async Task<Result<bool>> AddDutyAsync(DutyRosterDto dto)
{
if (string.IsNullOrWhiteSpace(dto?.PersonName))
{
return Result<bool>.Error("值班人不能为空");
}
try
{
var entity = dto.ToEntity();
entity.Id = 0;
entity.CreateTime = DateTime.Now;
await SqlSugarContext.DbContext.Insertable(entity).ExecuteCommandAsync();
return Result<bool>.Success(true);
}
catch (Exception ex)
{
return Result<bool>.Error("新增值班排班失败", ex);
}
}
/// <summary>
/// 修改排班
/// </summary>
public async Task<Result<bool>> UpdateDutyAsync(DutyRosterDto dto)
{
var entity = dto?.ToEntity();
if (entity == null || entity.Id <= 0)
{
return Result<bool>.Error("排班 Id 无效");
}
try
{
var rows = await SqlSugarContext.DbContext.Updateable(entity)
.IgnoreColumns(x => new { x.CreateTime, x.IsDel })
.ExecuteCommandAsync();
return Result<bool>.Success(rows > 0);
}
catch (Exception ex)
{
return Result<bool>.Error("修改值班排班失败", ex);
}
}
/// <summary>
/// 删除排班(软删除)
/// </summary>
public async Task<Result<bool>> DeleteDutyAsync(long id)
{
try
{
var rows = await SqlSugarContext.DbContext.Updateable<DutyRosterEntity>()
.SetColumns(x => x.IsDel == 1)
.Where(x => x.Id == id)
.ExecuteCommandAsync();
return Result<bool>.Success(rows > 0);
}
catch (Exception ex)
{
return Result<bool>.Error("删除值班排班失败", ex);
}
}
/// <summary>
/// 查询指定日期的值班人
/// </summary>
public async Task<DutyRosterEntity?> GetDutyByDateAsync(DateTime date)
{
return await SqlSugarContext.DbContext.Queryable<DutyRosterEntity>()
.Where(x => x.IsDel == 0 && x.DutyDate == date.Date)
.FirstAsync();
}
#endregion
#region 跨网发件箱(AlertBridge 对接)
/// <summary>
/// 写入推送日志(供 AlertNotifyWorker 直连推送后复用)
/// </summary>
public async Task WriteLogAsync(long alertId, long channelId, string content, bool success, string? errorMsg, int retryCount)
{
var channel = await SqlSugarContext.DbContext.Queryable<AlertNotifyChannelEntity>()
.Where(x => x.Id == channelId)
.FirstAsync();
if (channel == null) return;
await InsertLogAsync(alertId, channel, content,
success ? NotifyResult.Ok() : NotifyResult.Fail(errorMsg ?? "未知错误"), retryCount);
}
/// <summary>
/// 写入跨网发件箱(供 AlertNotifyWorker 的 Outbox 模式复用)
/// </summary>
public async Task EnqueueOutboxAsync(long alertId, long channelId, AlertNotifyMessage message, List<string> atMobiles)
{
var channel = await SqlSugarContext.DbContext.Queryable<AlertNotifyChannelEntity>()
.Where(x => x.Id == channelId && x.IsDel == 0)
.FirstAsync();
if (channel == null) return;
await InsertOutboxAsync(alertId, channel, message, atMobiles);
}
/// <summary>
/// 拉取待推送发件箱(先查后标记;单跳板机实例场景足够,多实例并发时可改用数据库锁)
/// </summary>
public async Task<Result<List<AlertOutboxItemDto>>> PullOutboxAsync(int batch)
{
try
{
if (batch <= 0 || batch > 100) batch = 20;
var pending = await SqlSugarContext.DbContext.Queryable<AlertOutboxEntity>()
.Where(x => x.IsDel == 0 && x.Status == OutboxStatusEnum.Pending)
.OrderBy(x => x.CreateTime, OrderByType.Asc)
.Take(batch)
.ToListAsync();
if (pending.Count == 0)
{
return Result<List<AlertOutboxItemDto>>.Success(new List<AlertOutboxItemDto>());
}
// 标记已拉取(仅更新仍为 Pending 的记录,防止重复拉取)
var ids = pending.Select(x => x.Id).ToList();
await SqlSugarContext.DbContext.Updateable<AlertOutboxEntity>()
.SetColumns(x => new AlertOutboxEntity { Status = OutboxStatusEnum.Pulled, PullTime = DateTime.Now })
.Where(x => ids.Contains(x.Id) && x.Status == OutboxStatusEnum.Pending)
.ExecuteCommandAsync();
return Result<List<AlertOutboxItemDto>>.Success(pending.ToDtoList());
}
catch (Exception ex)
{
return Result<List<AlertOutboxItemDto>>.Error("拉取发件箱失败", ex);
}
}
/// <summary>
/// 回写发件箱推送结果(失败累加重试次数;结果同步写推送日志便于页面追溯)
/// </summary>
public async Task<Result<bool>> ReportOutboxResultAsync(long id, OutboxResultDto dto)
{
try
{
var item = await SqlSugarContext.DbContext.Queryable<AlertOutboxEntity>()
.Where(x => x.Id == id && x.IsDel == 0)
.FirstAsync();
if (item == null)
{
return Result<bool>.Error("发件箱记录不存在");
}
var status = dto?.Success == true ? OutboxStatusEnum.Success : OutboxStatusEnum.Failed;
// 表达式树不支持 ?. 空传播,先提取局部变量
bool success = dto?.Success == true;
string? resultMsg = dto?.ResultMsg;
await SqlSugarContext.DbContext.Updateable<AlertOutboxEntity>()
.SetColumns(x => new AlertOutboxEntity
{
Status = status,
FinishTime = DateTime.Now,
ResultMsg = resultMsg,
RetryCount = success ? x.RetryCount : x.RetryCount + 1
})
.Where(x => x.Id == id)
.ExecuteCommandAsync();
// 推送日志(渠道信息回查,失败忽略)
var channel = await SqlSugarContext.DbContext.Queryable<AlertNotifyChannelEntity>()
.Where(x => x.Id == item.ChannelId)
.FirstAsync();
if (channel != null)
{
await InsertLogAsync(item.AlertId, channel,
$"[跨网中转] {(dto?.Success == true ? "推送成功" : "推送失败")}",
dto?.Success == true ? NotifyResult.Ok() : NotifyResult.Fail(dto?.ResultMsg ?? "未知错误"),
item.RetryCount);
}
return Result<bool>.Success(true);
}
catch (Exception ex)
{
return Result<bool>.Error("回写发件箱结果失败", ex);
}
}
/// <summary>
/// 写入发件箱(Worker 的 Outbox 模式与渠道测试推送共用;Payload 自包含全部推送要素)
/// </summary>
internal async Task InsertOutboxAsync(long alertId, AlertNotifyChannelEntity channel, AlertNotifyMessage message, List<string> atMobiles)
{
var payload = new AlertOutboxPayload
{
ChannelType = (int)channel.ChannelType,
WebhookUrl = channel.WebhookUrl,
Secret = channel.Secret,
Title = message.Title,
Content = message.Content,
AtMobiles = atMobiles
};
await SqlSugarContext.DbContext.Insertable(new AlertOutboxEntity
{
AlertId = alertId,
ChannelId = channel.Id,
Payload = System.Text.Json.JsonSerializer.Serialize(payload),
Status = OutboxStatusEnum.Pending,
CreateTime = DateTime.Now
}).ExecuteCommandAsync();
}
#endregion
}
}