OutsourceIssueMdpSyncService.cs 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421
  1. using Admin.NET.Plugin.AiDOP.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  3. using Admin.NET.Plugin.AiDOP.Infrastructure;
  4. namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  5. /// <summary>
  6. /// S5 委外发料单 数据中台只读同步转换服务(DOP 内部,head-detail 双 std 专表,独立于 S5MdpSyncTransformService 的 KPI 管线)。
  7. ///
  8. /// B-① 迁通用管线:源(NbrMaster/NbrDetail, Type='CA') 经执行器灌 mdp_stg_outsource_issue_pull(头/明细共表,source_table 区分)
  9. /// → transform 读 pull stg(MdpJsonSql 跨源类型兼容)→ mdp_std_outsource_issue(_detail)(typed,明细回填 std_head_id)。
  10. /// 维表 DepartmentMaster(d) 仍本库 LEFT JOIN(富化 department_name,业务 code 由事实表保留)。
  11. /// 注:canonical 执行器 stg 用 *_pull 命名(沿用 mdp_stg_ipqc_pull 惯例),与 legacy bespoke mdp_stg_outsource_issue(_detail) 区分;后者迁移后为死支路、不再消费、不动。
  12. ///
  13. /// 双源:本库实体 S5_OUTSOURCE_ISSUE_MASTER/_DETAIL(AIDOPDEV_MYSQL, status=1);
  14. /// dopdemorq 第二源 S5_OUTSOURCE_ISSUE_MASTER_SQLSERVER/_DETAIL_SQLSERVER(DOPDEMORQ_SQLSERVER, status=0 就位不启用)。
  15. ///
  16. /// 约束:
  17. /// - 只读源/贴源,仅写 mdp_stg_outsource_issue_pull / mdp_std_outsource_issue(_detail);绝不 INSERT/UPDATE/DELETE NbrMaster/NbrDetail。
  18. /// - 过滤 NbrMaster.Type='CA' AND IsActive=1;NbrDetail.Type='CA'。
  19. /// - status_desc 规则本批仅:Status='C' -> '关闭',其他 -> ''。department_name 经 Domain+Department 关联 DepartmentMaster.Descr。
  20. /// - 明细只转换已确认 9 列(ItemNum/ItemName/UM/QtyOrd/LocationFrom/LocationTo/Line/Status/Remark);发料数量/已发数/批次号 3 候选列本批后置。
  21. /// - tenant_id 从 raw_data 保留(本库源含该列 → std 租户不回归;SQL Server 源无此列 → 回落入站上下文 tenant)。
  22. /// - CA 为 0 行时转换成功完成、处理数为 0,不报错。切源 FULL Replace 头/明细各一次单事务、仅按 tenant、Pull 成功后才 destructive。
  23. /// </summary>
  24. public class OutsourceIssueMdpSyncService : ITransient
  25. {
  26. private const string JobCode = "S5_OUTSOURCE_ISSUE_MDP_SYNC";
  27. private const string InboundMasterEntityCode = "S5_OUTSOURCE_ISSUE_MASTER";
  28. private const string InboundDetailEntityCode = "S5_OUTSOURCE_ISSUE_DETAIL";
  29. private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
  30. private const string SqlServerMasterEntityCode = "S5_OUTSOURCE_ISSUE_MASTER_SQLSERVER";
  31. private const string SqlServerDetailEntityCode = "S5_OUTSOURCE_ISSUE_DETAIL_SQLSERVER";
  32. /// <summary>165 SQL Server 源无 tenant_id 列;切源时须经账套映射解析。</summary>
  33. private const string SourceZtid = "pbxfxp";
  34. private const string StdHeadTable = "mdp_std_outsource_issue";
  35. private const string StdDetailTable = "mdp_std_outsource_issue_detail";
  36. private readonly ISqlSugarClient _db;
  37. private readonly TransformRunLogFinalizer _runLogFinalizer;
  38. private readonly MdpSourcePullDispatcher _pullDispatcher;
  39. private readonly MdpNeutralSourceGate _neutralGate;
  40. public OutsourceIssueMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher, TransformRunLogFinalizer runLogFinalizer, MdpNeutralSourceGate neutralGate)
  41. {
  42. _db = db;
  43. _runLogFinalizer = runLogFinalizer;
  44. _pullDispatcher = pullDispatcher;
  45. _neutralGate = neutralGate;
  46. }
  47. /// <summary>全量:本地 DB 执行器灌 pull stg(头+明细) → 标准层头/明细(读 stg)。</summary>
  48. public async Task<OutsourceIssueMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  49. {
  50. cancellationToken.ThrowIfCancellationRequested();
  51. await EnsureTablesAsync();
  52. await EnsureStgTableAsync();
  53. var now = DateTime.Now;
  54. var batchId = $"S5_OUTSRC_ISSUE_FULL_{now:yyyyMMddHHmmss}";
  55. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  56. var result = new OutsourceIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  57. try
  58. {
  59. // 本库 NbrMaster 含 tenant_id;ctx tenantId=0 仅兜底,Writer 优先取源行值。
  60. var pullCtx = new MdpPullContext
  61. {
  62. TenantId = 0,
  63. FullRefresh = true,
  64. TaskCode = "S5_OUTSOURCE_ISSUE_INBOUND",
  65. BatchId = $"{batchId}_PULL"
  66. };
  67. var pull = await PopulateStgAsync(pullCtx, cancellationToken);
  68. result.HeadStageRows = pull.RowsWritten;
  69. result.DetailStageRows = 0;
  70. result.HeadStandardRows = await TransformHeadStandardAsync(batchId, now);
  71. result.DetailStandardRows = await TransformDetailStandardAsync(batchId, now);
  72. await MarkRunSuccessAsync(runLogId, now, result);
  73. return result;
  74. }
  75. catch (Exception ex)
  76. {
  77. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  78. if (!_runLogFinalizer.IsHostStopping)
  79. await MarkRunFailedAsync(runLogId, now, ex.Message);
  80. throw;
  81. }
  82. finally
  83. {
  84. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  85. }
  86. }
  87. /// <summary>双模式入站:执行器抽主/明细 → pull stg,再跑 CA 标准层头/明细(读 stg)。</summary>
  88. public async Task<OutsourceIssueInboundResult> RunInboundAsync(
  89. long tenantId = 0,
  90. bool fullRefresh = false,
  91. CancellationToken cancellationToken = default)
  92. {
  93. cancellationToken.ThrowIfCancellationRequested();
  94. await EnsureTablesAsync();
  95. await EnsureStgTableAsync();
  96. var now = DateTime.Now;
  97. // 入参 tenantId:本库源含 tenant_id 时可传 0(Writer 优先取源行)。
  98. var pullCtx = new MdpPullContext
  99. {
  100. TenantId = tenantId,
  101. FullRefresh = fullRefresh,
  102. TaskCode = "S5_OUTSOURCE_ISSUE_INBOUND",
  103. BatchId = $"S5_OUTSRC_ISSUE_IN_{now:yyyyMMddHHmmss}"
  104. };
  105. var pull = await PopulateStgAsync(pullCtx, cancellationToken);
  106. var batchId = $"S5_OUTSRC_ISSUE_STD_{now:yyyyMMddHHmmss}";
  107. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  108. var result = new OutsourceIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  109. try
  110. {
  111. result.HeadStandardRows = await TransformHeadStandardAsync(batchId, now);
  112. result.DetailStandardRows = await TransformDetailStandardAsync(batchId, now);
  113. await MarkRunSuccessAsync(runLogId, now, result);
  114. }
  115. catch (Exception ex)
  116. {
  117. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  118. if (!_runLogFinalizer.IsHostStopping)
  119. await MarkRunFailedAsync(runLogId, now, ex.Message);
  120. throw;
  121. }
  122. finally
  123. {
  124. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  125. }
  126. return new OutsourceIssueInboundResult
  127. {
  128. PullBatchId = pullCtx.BatchId,
  129. RowsPulled = pull.RowsPulled,
  130. RowsWrittenStg = pull.RowsWritten,
  131. NewCursor = pull.NewCursor,
  132. TransformBatchId = batchId,
  133. StdHeadRows = result.HeadStandardRows,
  134. StdDetailRows = result.DetailStandardRows
  135. };
  136. }
  137. /// <summary>
  138. /// 双源切换 FULL Replace:从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 pull stg,
  139. /// 成功后头/明细各一次单事务 FULL 重建(按 tenant 精确隔离),消除旧源独有业务键残留。
  140. /// 顺序:先 head(写 std 头,供明细回填 std_head_id),再 detail。两次独立事务,Pull 成功后才 destructive。
  141. /// SQLSERVER 实体默认 status=0(就位不启用);dopdemorq 六表当前空,Phase 1 不实际执行 destructive 切换。
  142. /// </summary>
  143. public async Task<OutsourceIssueInboundResult> RunSourceSwitchFullAsync(
  144. string sourceCode = SqlServerSourceCode,
  145. string masterEntityCode = SqlServerMasterEntityCode,
  146. string detailEntityCode = SqlServerDetailEntityCode,
  147. long tenantId = 0,
  148. CancellationToken cancellationToken = default)
  149. {
  150. cancellationToken.ThrowIfCancellationRequested();
  151. await EnsureTablesAsync();
  152. await EnsureStgTableAsync();
  153. // 165 SQL Server 源无 tenant_id;显式 tenantId≤0 时经账套映射解析。
  154. tenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  155. var now = DateTime.Now;
  156. var pullCtx = new MdpPullContext
  157. {
  158. TenantId = tenantId,
  159. FullRefresh = true,
  160. TaskCode = "S5_OUTSOURCE_ISSUE_INBOUND",
  161. BatchId = $"S5_OUTSRC_ISSUE_SW_{now:yyyyMMddHHmmss}"
  162. };
  163. var master = await _pullDispatcher.PullAllByEntityCodeAsync(masterEntityCode, pullCtx, cancellationToken);
  164. var detail = await _pullDispatcher.PullAllByEntityCodeAsync(detailEntityCode, pullCtx, cancellationToken);
  165. var batchId = $"S5_OUTSRC_ISSUE_SWSTD_{now:yyyyMMddHHmmss}";
  166. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  167. var result = new OutsourceIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  168. try
  169. {
  170. result.HeadStandardRows = await MdpStdFullReplace.ReplaceAsync(
  171. _db, StdHeadTable, tenantId, extraWhere: null,
  172. insertScopedAsync: () => TransformHeadStandardAsync(batchId, now, sourceCode),
  173. cancellationToken);
  174. result.DetailStandardRows = await MdpStdFullReplace.ReplaceAsync(
  175. _db, StdDetailTable, tenantId, extraWhere: null,
  176. insertScopedAsync: () => TransformDetailStandardAsync(batchId, now, sourceCode),
  177. cancellationToken);
  178. await MarkRunSuccessAsync(runLogId, now, result);
  179. }
  180. catch (Exception ex)
  181. {
  182. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  183. if (!_runLogFinalizer.IsHostStopping)
  184. await MarkRunFailedAsync(runLogId, now, ex.Message);
  185. throw;
  186. }
  187. finally
  188. {
  189. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  190. }
  191. return new OutsourceIssueInboundResult
  192. {
  193. PullBatchId = pullCtx.BatchId,
  194. RowsPulled = master.RowsPulled + detail.RowsPulled,
  195. RowsWrittenStg = master.RowsWritten + detail.RowsWritten,
  196. NewCursor = detail.NewCursor ?? master.NewCursor,
  197. TransformBatchId = batchId,
  198. StdHeadRows = result.HeadStandardRows,
  199. StdDetailRows = result.DetailStandardRows
  200. };
  201. }
  202. private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
  203. MdpPullContext pullCtx, CancellationToken cancellationToken)
  204. {
  205. var master = await _pullDispatcher.PullByEntityCodeAsync(InboundMasterEntityCode, pullCtx, cancellationToken);
  206. var detail = await _pullDispatcher.PullByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
  207. return (
  208. master.RowsPulled + detail.RowsPulled,
  209. master.RowsWritten + detail.RowsWritten,
  210. detail.NewCursor ?? master.NewCursor);
  211. }
  212. private async Task EnsureStgTableAsync()
  213. {
  214. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_outsource_issue_pull"));
  215. }
  216. /// <summary>防御式建表(与 UpdateScripts/1.0.206.sql 同构,幂等)。std 头/明细两专表。</summary>
  217. private async Task EnsureTablesAsync()
  218. {
  219. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_outsource_issue"));
  220. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_outsource_issue_detail"));
  221. }
  222. /// <summary>标准化头:pull stg(NbrMaster, CA) + DepartmentMaster(本库) → mdp_std_outsource_issue。返回处理头行数。</summary>
  223. private async Task<int> TransformHeadStandardAsync(string batchId, DateTime now, string? sourceSystem = null)
  224. {
  225. var gateSource = sourceSystem ?? MdpSourceIdentity.Native;
  226. if (!await _neutralGate.AllowsAsync(0, "OUTSOURCE", gateSource))
  227. return 0;
  228. var srcClause = sourceSystem == null ? "" : " AND m.source_system=@Src";
  229. var mTenant = $"COALESCE({MdpJsonSql.Int("m", "tenant_id")}, NULLIF(m.tenant_id, 0))";
  230. var statusExpr = MdpJsonSql.Str("m", "Status");
  231. var headWhere =
  232. $"""
  233. m.source_table='NbrMaster'
  234. AND {MdpJsonSql.Str("m", "Type")}='CA'
  235. AND {MdpJsonSql.BoolTrue("m", "IsActive")}{srcClause}
  236. AND {MdpJsonSql.TenantGuard(mTenant)}
  237. """;
  238. var countPars = new List<SugarParameter>();
  239. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  240. var rows = await _db.Ado.GetIntAsync(
  241. $"SELECT COUNT(1) FROM mdp_stg_outsource_issue_pull m WHERE {headWhere}", countPars);
  242. var insertSql =
  243. $"""
  244. INSERT INTO mdp_std_outsource_issue
  245. (tenant_id, factory_id, source_system, bill_no, issue_date, outsource_no, work_order,
  246. department_code, department_name, issuer, status, status_desc, remark, create_user,
  247. source_create_time, eff_date, source_biz_key, sync_batch_id, sync_time)
  248. SELECT
  249. {mTenant}, 1, {MdpSourceIdentity.Resolve("m.source_table", "m.source_system")},
  250. {MdpJsonSql.Str("m", "Nbr")}, {MdpJsonSql.DateTimeSec("m", "Date")}, {MdpJsonSql.Str("m", "Address")}, {MdpJsonSql.Str("m", "WorkOrd")},
  251. {MdpJsonSql.Str("m", "Department")}, d.Descr, {MdpJsonSql.Str("m", "User1")}, {statusExpr},
  252. CASE WHEN {statusExpr}='C' THEN '关闭' ELSE '' END,
  253. {MdpJsonSql.Str("m", "Remark")}, {MdpJsonSql.Str("m", "CreateUser")}, {MdpJsonSql.DateTimeSec("m", "CreateTime")}, {MdpJsonSql.DateTimeSec("m", "EffDate")},
  254. {MdpJsonSql.Str("m", "Nbr")}, @BatchId, @Now
  255. FROM mdp_stg_outsource_issue_pull m
  256. LEFT JOIN DepartmentMaster d
  257. ON d.Domain = {MdpJsonSql.Str("m", "Domain")} AND d.Department = {MdpJsonSql.Str("m", "Department")}
  258. WHERE {headWhere}
  259. ON DUPLICATE KEY UPDATE
  260. issue_date=VALUES(issue_date), outsource_no=VALUES(outsource_no), work_order=VALUES(work_order),
  261. department_code=VALUES(department_code), department_name=VALUES(department_name), issuer=VALUES(issuer),
  262. status=VALUES(status), status_desc=VALUES(status_desc), remark=VALUES(remark), create_user=VALUES(create_user),
  263. source_create_time=VALUES(source_create_time), eff_date=VALUES(eff_date),
  264. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  265. """;
  266. var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
  267. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  268. await MdpSchemaAligner.ExecuteAsync(_db, insertSql, insPars);
  269. return rows;
  270. }
  271. /// <summary>标准化明细:pull stg(NbrDetail, CA) self-join stg(NbrMaster) by RecID↔NbrRecID → mdp_std_outsource_issue_detail(仅 9 列;回填 std_head_id)。</summary>
  272. private async Task<int> TransformDetailStandardAsync(string batchId, DateTime now, string? sourceSystem = null)
  273. {
  274. var gateSource = sourceSystem ?? MdpSourceIdentity.Native;
  275. if (!await _neutralGate.AllowsAsync(0, "OUTSOURCE", gateSource))
  276. return 0;
  277. var srcClause = sourceSystem == null ? "" : " AND m.source_system=@Src AND d.source_system=@Src";
  278. var dTenant = $"COALESCE({MdpJsonSql.Int("d", "tenant_id")}, NULLIF(d.tenant_id, 0))";
  279. var fromJoinWhere =
  280. $"""
  281. FROM mdp_stg_outsource_issue_pull d
  282. JOIN mdp_stg_outsource_issue_pull m
  283. ON m.source_table='NbrMaster'
  284. AND {MdpJsonSql.Int("m", "RecID")} = {MdpJsonSql.Int("d", "NbrRecID")}
  285. AND {MdpJsonSql.Str("m", "Type")}='CA'
  286. AND {MdpJsonSql.BoolTrue("m", "IsActive")}
  287. LEFT JOIN mdp_std_outsource_issue h
  288. ON h.tenant_id = {dTenant} AND h.bill_no = {MdpJsonSql.Str("d", "Nbr")}
  289. WHERE d.source_table='NbrDetail' AND {MdpJsonSql.Str("d", "Type")}='CA'{srcClause}
  290. AND {MdpJsonSql.TenantGuard(dTenant)}
  291. """;
  292. var countPars = new List<SugarParameter>();
  293. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  294. var rows = await _db.Ado.GetIntAsync($"SELECT COUNT(1) {fromJoinWhere}", countPars);
  295. var insertSql =
  296. $"""
  297. INSERT INTO mdp_std_outsource_issue_detail
  298. (tenant_id, std_head_id, bill_no, line, item_num, item_name, um, qty_ord,
  299. location_from, location_to, status, remark, source_biz_key, sync_batch_id, sync_time)
  300. SELECT
  301. {dTenant}, h.id, {MdpJsonSql.Str("d", "Nbr")}, {MdpJsonSql.Int("d", "Line")}, {MdpJsonSql.Str("d", "ItemNum")}, {MdpJsonSql.Str("d", "ItemName")}, {MdpJsonSql.Str("d", "UM")}, {MdpJsonSql.Dec("d", "QtyOrd", 18, 6)},
  302. {MdpJsonSql.Str("d", "LocationFrom")}, {MdpJsonSql.Str("d", "LocationTo")}, {MdpJsonSql.Str("d", "Status")}, {MdpJsonSql.Str("d", "Remark")},
  303. CONCAT(IFNULL({MdpJsonSql.Str("d", "Nbr")},''), '#', IFNULL({MdpJsonSql.Int("d", "Line")},0)), @BatchId, @Now
  304. {fromJoinWhere}
  305. ON DUPLICATE KEY UPDATE
  306. std_head_id=VALUES(std_head_id), line=VALUES(line), item_num=VALUES(item_num), item_name=VALUES(item_name),
  307. um=VALUES(um), qty_ord=VALUES(qty_ord), location_from=VALUES(location_from), location_to=VALUES(location_to),
  308. status=VALUES(status), remark=VALUES(remark), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  309. update_time=CURRENT_TIMESTAMP
  310. """;
  311. var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
  312. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  313. await MdpSchemaAligner.ExecuteAsync(_db, insertSql, insPars);
  314. return rows;
  315. }
  316. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  317. {
  318. await MdpSchemaAligner.ExecuteAsync(_db,
  319. """
  320. INSERT INTO mdp_transform_run_log
  321. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  322. VALUES (0, @JobCode, 'S5委外发料单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  323. """,
  324. new SugarParameter("@JobCode", JobCode),
  325. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  326. new SugarParameter("@BatchId", batchId),
  327. new SugarParameter("@StartTime", startedAt));
  328. return await _db.Ado.GetLongAsync(
  329. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  330. new List<SugarParameter> { new("@BatchId", batchId) });
  331. }
  332. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, OutsourceIssueMdpSyncResult result)
  333. {
  334. var finishedAt = DateTime.Now;
  335. await MdpSchemaAligner.ExecuteAsync(_db,
  336. """
  337. UPDATE mdp_transform_run_log
  338. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  339. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  340. WHERE id=@Id
  341. """,
  342. new SugarParameter("@EndTime", finishedAt),
  343. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  344. new SugarParameter("@StageRows", result.HeadStageRows + result.DetailStageRows),
  345. new SugarParameter("@StandardRows", result.HeadStandardRows + result.DetailStandardRows),
  346. new SugarParameter("@Id", runLogId));
  347. }
  348. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  349. {
  350. var finishedAt = DateTime.Now;
  351. await MdpSchemaAligner.ExecuteAsync(_db,
  352. """
  353. UPDATE mdp_transform_run_log
  354. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  355. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  356. WHERE id=@Id
  357. """,
  358. new SugarParameter("@EndTime", finishedAt),
  359. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  360. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  361. new SugarParameter("@Id", runLogId));
  362. }
  363. private static string NormalizeTriggerType(string? triggerType)
  364. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  365. }
  366. /// <summary>委外发料单 MDP 同步转换结果。</summary>
  367. public sealed class OutsourceIssueMdpSyncResult
  368. {
  369. public long RunLogId { get; set; }
  370. public string BatchId { get; set; } = string.Empty;
  371. public int HeadStageRows { get; set; }
  372. public int DetailStageRows { get; set; }
  373. public int HeadStandardRows { get; set; }
  374. public int DetailStandardRows { get; set; }
  375. }
  376. /// <summary>S5 委外发料双模式入站结果。</summary>
  377. public sealed class OutsourceIssueInboundResult
  378. {
  379. public string PullBatchId { get; set; } = string.Empty;
  380. public int RowsPulled { get; set; }
  381. public int RowsWrittenStg { get; set; }
  382. public string? NewCursor { get; set; }
  383. public string TransformBatchId { get; set; } = string.Empty;
  384. public int StdHeadRows { get; set; }
  385. public int StdDetailRows { get; set; }
  386. }