S3MdpMonitorService.cs 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144
  1. using Admin.NET.Plugin.AiDOP.Order;
  2. namespace Admin.NET.Plugin.AiDOP.Supply;
  3. /// <summary>
  4. /// S3 MDP 运行监控。
  5. ///
  6. /// 【租户安全边界】租户一律经 <see cref="AidopTenantScope.ResolveOrThrow"/> 从认证后 JWT 解析,
  7. /// 无有效租户即拒绝,不读前端 tenantId、无默认回退(原实现:类级 <c>[AllowAnonymous]</c> +
  8. /// <c>_userManager.TenantId &gt; 0</c> 跳过分支,匿名时退化为 <c>tenant_id = 0</c> 平台行或不过滤)。
  9. /// 注:租户行与 <c>tenant_id = 0</c> 平台行的可见性语义(<c>BuildMdpRunLogTenantWhere</c>)保持不变。
  10. /// </summary>
  11. [ApiDescriptionSettings(Order = 320, Description = "S3 MDP运行监控")]
  12. [Route("api/Supply")]
  13. [NonUnify]
  14. public class S3MdpMonitorService : IDynamicApiController, ITransient
  15. {
  16. private readonly ISqlSugarClient _db;
  17. private readonly UserManager _userManager;
  18. public S3MdpMonitorService(ISqlSugarClient db, UserManager userManager)
  19. {
  20. _db = db;
  21. _userManager = userManager;
  22. }
  23. [DisplayName("S3 MDP最近运行状态")]
  24. [HttpGet("s3-mdp-monitor/latest")]
  25. public async Task<object> GetLatest()
  26. {
  27. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  28. var tenantWhere = $" AND {MdpMonitorService.BuildMdpRunLogTenantWhere(tenantId)}";
  29. var pars = new List<SugarParameter> { new("@TenantId", tenantId) };
  30. return await _db.Ado.SqlQuerySingleAsync<S3MdpRunLogRow>(
  31. $"{SelectColumnsSql()} FROM mdp_transform_run_log WHERE job_code='S3_MDP_SYNC_TRANSFORM'{tenantWhere} ORDER BY start_time DESC, id DESC LIMIT 1",
  32. pars)
  33. ?? new S3MdpRunLogRow();
  34. }
  35. [DisplayName("S3 MDP运行日志列表")]
  36. [HttpGet("s3-mdp-monitor/list")]
  37. public async Task<object> GetList([FromQuery] S3MdpMonitorListInput input)
  38. {
  39. var page = input.Page <= 0 ? 1 : input.Page;
  40. var pageSize = input.PageSize <= 0 ? 10 : input.PageSize;
  41. var offset = (page - 1) * pageSize;
  42. var where = new List<string> { "job_code='S3_MDP_SYNC_TRANSFORM'" };
  43. var pars = new List<SugarParameter>();
  44. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  45. where.Add(MdpMonitorService.BuildMdpRunLogTenantWhere(tenantId));
  46. pars.Add(new SugarParameter("@TenantId", tenantId));
  47. if (!string.IsNullOrWhiteSpace(input.BatchId))
  48. {
  49. where.Add("batch_id LIKE @BatchId");
  50. pars.Add(new SugarParameter("@BatchId", $"%{input.BatchId.Trim()}%"));
  51. }
  52. if (!string.IsNullOrWhiteSpace(input.Status))
  53. {
  54. where.Add("status=@Status");
  55. pars.Add(new SugarParameter("@Status", input.Status.Trim().ToUpperInvariant()));
  56. }
  57. if (input.StartTime.HasValue)
  58. {
  59. where.Add("start_time >= @StartTime");
  60. pars.Add(new SugarParameter("@StartTime", input.StartTime.Value));
  61. }
  62. if (input.EndTime.HasValue)
  63. {
  64. where.Add("start_time <= @EndTime");
  65. pars.Add(new SugarParameter("@EndTime", input.EndTime.Value));
  66. }
  67. var whereSql = string.Join(" AND ", where);
  68. var total = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM mdp_transform_run_log WHERE {whereSql}", pars);
  69. var list = await _db.Ado.SqlQueryAsync<S3MdpRunLogRow>(
  70. $"""
  71. {SelectColumnsSql()}
  72. FROM mdp_transform_run_log
  73. WHERE {whereSql}
  74. ORDER BY start_time DESC, id DESC
  75. LIMIT {pageSize} OFFSET {offset}
  76. """,
  77. pars);
  78. return new { total, page, pageSize, list };
  79. }
  80. [DisplayName("S3 MDP运行日志详情")]
  81. [HttpGet("s3-mdp-monitor/detail/{id}")]
  82. public async Task<object> GetDetail(long id)
  83. {
  84. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  85. var pars = new List<SugarParameter> { new("@Id", id), new("@TenantId", tenantId) };
  86. var row = await _db.Ado.SqlQuerySingleAsync<S3MdpRunLogRow>(
  87. $"{SelectColumnsSql()} FROM mdp_transform_run_log WHERE id=@Id AND tenant_id=@TenantId LIMIT 1",
  88. pars);
  89. return row ?? throw Oops.Oh("运行日志不存在");
  90. }
  91. private static string SelectColumnsSql()
  92. {
  93. return """
  94. SELECT id AS Id, tenant_id AS TenantId, job_code AS JobCode, job_name AS JobName, trigger_type AS TriggerType,
  95. batch_id AS BatchId, status AS Status, start_time AS StartTime, end_time AS EndTime, duration_ms AS DurationMs,
  96. stage_rows AS StageRows, standard_rows AS StandardRows, dwd_rows AS DwdRows,
  97. error_message AS ErrorMessage, summary_json AS SummaryJson, create_time AS CreateTime, update_time AS UpdateTime
  98. """;
  99. }
  100. }
  101. public sealed class S3MdpMonitorListInput
  102. {
  103. public string? BatchId { get; set; }
  104. public string? Status { get; set; }
  105. public DateTime? StartTime { get; set; }
  106. public DateTime? EndTime { get; set; }
  107. public int Page { get; set; } = 1;
  108. public int PageSize { get; set; } = 10;
  109. }
  110. public sealed class S3MdpRunLogRow
  111. {
  112. public long Id { get; set; }
  113. public long TenantId { get; set; }
  114. public string? JobCode { get; set; }
  115. public string? JobName { get; set; }
  116. public string? TriggerType { get; set; }
  117. public string? BatchId { get; set; }
  118. public string? Status { get; set; }
  119. public DateTime? StartTime { get; set; }
  120. public DateTime? EndTime { get; set; }
  121. public int? DurationMs { get; set; }
  122. public int? StageRows { get; set; }
  123. public int? StandardRows { get; set; }
  124. public int? DwdRows { get; set; }
  125. public string? ErrorMessage { get; set; }
  126. public string? SummaryJson { get; set; }
  127. public DateTime? CreateTime { get; set; }
  128. public DateTime? UpdateTime { get; set; }
  129. }