ReportWorkMdpSyncService.cs 10.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210
  1. using System.Text.Json;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform;
  3. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  4. using Admin.NET.Plugin.AiDOP.Infrastructure;
  5. namespace Admin.NET.Plugin.AiDOP.Manufacturing;
  6. /// <summary>
  7. /// S6 报工双模式入站:执行器 → mdp_stg_s6_report → mdp_std_s6_report。
  8. /// 源优先 T8 <c>Cj_Bg_Head_Rep</c>(实体 S6_REPORT);API 对偶 S6_REPORT_API。
  9. /// </summary>
  10. public class ReportWorkMdpSyncService : ITransient
  11. {
  12. private const string InboundEntityCode = "S6_REPORT";
  13. /// <summary>T8 报工源无 tenant_id;须经 <see cref="AidopSourceTenantMap"/> 解析。</summary>
  14. private const string SourceZtid = "pbxfxp";
  15. private readonly ISqlSugarClient _db;
  16. private readonly MdpSourcePullDispatcher _pullDispatcher;
  17. private readonly MdpNeutralSourceGate _neutralGate;
  18. public ReportWorkMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher, MdpNeutralSourceGate neutralGate)
  19. {
  20. _db = db;
  21. _pullDispatcher = pullDispatcher;
  22. _neutralGate = neutralGate;
  23. }
  24. public async Task<ReportWorkInboundResult> RunInboundAsync(
  25. long tenantId = 0,
  26. bool fullRefresh = false,
  27. string? entityCode = null,
  28. CancellationToken cancellationToken = default)
  29. {
  30. cancellationToken.ThrowIfCancellationRequested();
  31. await EnsureTablesAsync();
  32. if (tenantId > 0 && await _neutralGate.AllowsAsync(tenantId, "S6_REPORT", "DOPDEMORQ_SQLSERVER"))
  33. return await Run165Async(tenantId, fullRefresh, cancellationToken);
  34. // T8 源表无 tenant_id 列;显式 tenantId≤0 时经账套映射解析真实租户。
  35. var resolvedTenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  36. var code = string.IsNullOrWhiteSpace(entityCode) ? InboundEntityCode : entityCode.Trim();
  37. var now = DateTime.Now;
  38. var pullCtx = new MdpPullContext
  39. {
  40. TenantId = resolvedTenantId,
  41. FullRefresh = fullRefresh,
  42. TaskCode = "S6_REPORT_INBOUND",
  43. BatchId = $"S6_RPT_IN_{now:yyyyMMddHHmmss}"
  44. };
  45. var pull = await _pullDispatcher.PullByEntityCodeAsync(code, pullCtx, cancellationToken);
  46. var transformBatch = $"{pullCtx.BatchId}_STD";
  47. var stdRows = await TransformStandardAsync(resolvedTenantId, transformBatch, now);
  48. return new ReportWorkInboundResult
  49. {
  50. PullBatchId = pullCtx.BatchId,
  51. RowsPulled = pull.RowsWritten,
  52. RowsWrittenStg = pull.RowsWritten,
  53. TransformBatchId = transformBatch,
  54. StandardRows = stdRows
  55. };
  56. }
  57. private async Task<int> TransformStandardAsync(long tenantId, string batchId, DateTime now)
  58. {
  59. if (!await _neutralGate.AllowsAsync(tenantId, "S6_REPORT", "T8"))
  60. return 0;
  61. // 将 PENDING stg 投影到 std(JSON 字段兼容大小写)
  62. var sTenant = MdpJsonSql.TenantFromStg("s", "@tid");
  63. var sql = $"""
  64. INSERT INTO mdp_std_s6_report
  65. (tenant_id, factory_id, source_system, work_order_no, report_date, report_qty, ztid,
  66. source_row_id, source_biz_key, sync_batch_id, sync_time)
  67. SELECT
  68. {sTenant},
  69. 1,
  70. IFNULL(NULLIF(s.source_system,''), 'T8'),
  71. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.noid')), 'null'),
  72. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.NOID')), 'null'), s.source_biz_key)),
  73. COALESCE(
  74. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.kgdate')), 'null'), ''), '%Y-%m-%d %H:%i:%s'),
  75. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.KGDATE')), 'null'), ''), '%Y-%m-%d %H:%i:%s'),
  76. NULL),
  77. CAST(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.sl')), 'null'), '') AS DECIMAL(18,6)),
  78. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.ztid')), 'null'),
  79. IFNULL(NULLIF(s.source_row_id,''), s.source_biz_key),
  80. s.source_biz_key,
  81. @batch,
  82. @now
  83. FROM mdp_stg_s6_report s
  84. WHERE s.process_status = 'PENDING'
  85. AND s.source_biz_key IS NOT NULL
  86. AND s.source_biz_key <> ''
  87. AND {MdpJsonSql.TenantGuard(sTenant)}
  88. ON DUPLICATE KEY UPDATE
  89. work_order_no = VALUES(work_order_no),
  90. report_date = VALUES(report_date),
  91. report_qty = VALUES(report_qty),
  92. ztid = VALUES(ztid),
  93. source_row_id = VALUES(source_row_id),
  94. sync_batch_id = VALUES(sync_batch_id),
  95. sync_time = VALUES(sync_time),
  96. update_time = CURRENT_TIMESTAMP
  97. """;
  98. var affected = await MdpSchemaAligner.ExecuteAsync(_db, sql,
  99. new SugarParameter("@tid", tenantId),
  100. new SugarParameter("@batch", batchId),
  101. new SugarParameter("@now", now));
  102. await MdpSchemaAligner.ExecuteAsync(_db,
  103. tenantId > 0
  104. ? "UPDATE mdp_stg_s6_report SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING' AND tenant_id = @tid"
  105. : "UPDATE mdp_stg_s6_report SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING'",
  106. tenantId > 0 ? new SugarParameter("@tid", tenantId) : null!);
  107. return affected;
  108. }
  109. private async Task<ReportWorkInboundResult> Run165Async(long tenantId, bool fullRefresh, CancellationToken cancellationToken)
  110. {
  111. var now = DateTime.Now;
  112. var pullCtx = new MdpPullContext
  113. {
  114. TenantId = tenantId,
  115. FullRefresh = fullRefresh,
  116. TaskCode = "S6_EMP_WORK_ORD_TRANS",
  117. BatchId = $"S6_165_RPT_{now:yyyyMMddHHmmss}"
  118. };
  119. var pull = await _pullDispatcher.PullByEntityCodeAsync("S6_EMP_WORK_ORD_TRANS_SQLSERVER", pullCtx, cancellationToken);
  120. await EnsureStartWorkColumnAsync();
  121. var batch = $"{pullCtx.BatchId}_STD";
  122. var stdRows = await _db.Ado.ExecuteCommandAsync(
  123. """
  124. INSERT INTO mdp_std_s6_report
  125. (tenant_id, factory_id, source_system, written_by, work_order_no, report_date, start_work_date, report_qty, ztid,
  126. source_row_id, source_biz_key, sync_batch_id, sync_time)
  127. SELECT
  128. @tid, 1, 'DOPDEMORQ_SQLSERVER', 'DB_SYNC',
  129. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrd')),'null'),
  130. COALESCE(
  131. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  132. STR_TO_DATE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdDate')),'null'), '%Y-%m-%d')),
  133. COALESCE(
  134. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.StartDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  135. STR_TO_DATE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.StartDate')),'null'), '%Y-%m-%d')),
  136. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CompQty')),'null') AS DECIMAL(18,6)),
  137. '',
  138. IFNULL(NULLIF(s.source_row_id,''), s.source_biz_key),
  139. s.source_biz_key, @batch, @now
  140. FROM mdp_stg_s6_report s
  141. WHERE s.tenant_id=@tid AND s.process_status='PENDING'
  142. AND s.source_system='DOPDEMORQ_SQLSERVER'
  143. ON DUPLICATE KEY UPDATE
  144. report_qty=VALUES(report_qty), report_date=VALUES(report_date),
  145. start_work_date=VALUES(start_work_date), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  146. """,
  147. new SugarParameter("@tid", tenantId),
  148. new SugarParameter("@batch", batch),
  149. new SugarParameter("@now", now));
  150. await _db.Ado.ExecuteCommandAsync(
  151. """
  152. UPDATE mdp_std_work_order_line w
  153. JOIN (
  154. SELECT work_order_no, MIN(start_work_date) d
  155. FROM mdp_std_s6_report
  156. WHERE tenant_id=@tid AND source_system='DOPDEMORQ_SQLSERVER' AND start_work_date IS NOT NULL
  157. GROUP BY work_order_no
  158. ) r ON r.work_order_no=w.order_no
  159. SET w.release_time=r.d
  160. WHERE w.tenant_id=@tid AND w.doc_type='PROD_TASK' AND w.source_system='AIDOP_NATIVE'
  161. """,
  162. new SugarParameter("@tid", tenantId));
  163. await _db.Ado.ExecuteCommandAsync(
  164. "UPDATE mdp_stg_s6_report SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING' AND tenant_id=@tid AND source_system='DOPDEMORQ_SQLSERVER'",
  165. new SugarParameter("@tid", tenantId));
  166. return new ReportWorkInboundResult
  167. {
  168. PullBatchId = pullCtx.BatchId,
  169. RowsPulled = pull.RowsWritten,
  170. RowsWrittenStg = pull.RowsWritten,
  171. TransformBatchId = batch,
  172. StandardRows = stdRows
  173. };
  174. }
  175. private async Task EnsureStartWorkColumnAsync()
  176. {
  177. var exists = await _db.Ado.GetIntAsync(
  178. "SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='mdp_std_s6_report' AND COLUMN_NAME='start_work_date'");
  179. if (exists == 0)
  180. await _db.Ado.ExecuteCommandAsync("ALTER TABLE mdp_std_s6_report ADD COLUMN start_work_date datetime NULL");
  181. await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_s6_report");
  182. }
  183. private async Task EnsureTablesAsync()
  184. {
  185. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_s6_report"));
  186. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_s6_report"));
  187. }
  188. }
  189. public class ReportWorkInboundResult
  190. {
  191. public string PullBatchId { get; set; } = string.Empty;
  192. public int RowsPulled { get; set; }
  193. public int RowsWrittenStg { get; set; }
  194. public string TransformBatchId { get; set; } = string.Empty;
  195. public int StandardRows { get; set; }
  196. }