ProductionReceiptMdpSyncService.cs 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405
  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 内部,扁平单表,独立于 S5MdpSyncTransformService 的 KPI 管线)。
  7. ///
  8. /// 源:aidopdev NbrMaster(n) + NbrDetail(d),业务类型 Type='WOI'(生产入库);
  9. /// 扁平 LIST join:n.Domain=d.Domain AND n.Nbr=d.Nbr(一行=一条 NbrDetail 明细)。
  10. /// 维表 DepartmentMaster(p) / ItemMaster(i) / LocationMaster(lt,lf) 全 LEFT JOIN。
  11. /// LocationMaster 无 Domain 列(DOP 重建:domain_code/location/descr),按 tenant_id+location join。
  12. /// 链路(Phase 1):执行器抽 NbrMaster/NbrDetail → mdp_stg_production_receipt → mdp_std_production_receipt。
  13. /// 双模式入站:mdp_entity=S7_NBR_MASTER/S7_NBR_DETAIL → Writer → stg;标准层读 stg;维表仍直连本库。
  14. ///
  15. /// 约束:
  16. /// - 只读源/贴源,仅写 mdp_stg_production_receipt / mdp_std_production_receipt;绝不写 NbrMaster/NbrDetail。
  17. /// - 过滤 Type='WOI' AND IsActive=1。
  18. /// - WOI 为 0 行时转换成功完成、处理数为 0,不报错。
  19. /// </summary>
  20. public class ProductionReceiptMdpSyncService : ITransient
  21. {
  22. private const string JobCode = "S5_PRODUCTION_RECEIPT_MDP_SYNC";
  23. private const string InboundEntityCode = "S7_NBR_MASTER";
  24. private const string InboundDetailEntityCode = "S7_NBR_DETAIL";
  25. // 双源(dopdemorq SQL Server)Phase 1:沿用现有 S7_NBR_* 命名语义(业务=生产入库),加 _SQLSERVER 第二源实体。实体默认 status=0(就位不启用)。
  26. private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
  27. private const string SqlServerMasterEntityCode = "S7_NBR_MASTER_SQLSERVER";
  28. private const string SqlServerDetailEntityCode = "S7_NBR_DETAIL_SQLSERVER";
  29. /// <summary>165 SQL Server 源无 tenant_id 列;切源时须经账套映射解析。</summary>
  30. private const string SourceZtid = "pbxfxp";
  31. private readonly ISqlSugarClient _db;
  32. private readonly TransformRunLogFinalizer _runLogFinalizer;
  33. private readonly MdpSourcePullDispatcher _pullDispatcher;
  34. private readonly MdpNeutralSourceGate _neutralGate;
  35. public ProductionReceiptMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher, TransformRunLogFinalizer runLogFinalizer, MdpNeutralSourceGate neutralGate)
  36. {
  37. _db = db;
  38. _runLogFinalizer = runLogFinalizer;
  39. _pullDispatcher = pullDispatcher;
  40. _neutralGate = neutralGate;
  41. }
  42. /// <summary>全量:本地 DB 执行器灌 stg → 标准层(读 stg)。</summary>
  43. public async Task<ProductionReceiptMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  44. {
  45. cancellationToken.ThrowIfCancellationRequested();
  46. await EnsureTablesAsync();
  47. await EnsureStgTableAsync();
  48. var now = DateTime.Now;
  49. var batchId = $"S5_PROD_RCPT_FULL_{now:yyyyMMddHHmmss}";
  50. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  51. var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  52. try
  53. {
  54. // 本库 NbrMaster 含 tenant_id;ctx tenantId=0 仅兜底,Writer 优先取源行值。
  55. var pullCtx = new MdpPullContext
  56. {
  57. TenantId = 0,
  58. FullRefresh = true,
  59. TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND",
  60. BatchId = $"{batchId}_PULL"
  61. };
  62. await PopulateStgAsync(pullCtx, cancellationToken);
  63. result.Rows = await TransformStandardAsync(batchId, now);
  64. await MarkRunSuccessAsync(runLogId, now, result);
  65. return result;
  66. }
  67. catch (Exception ex)
  68. {
  69. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  70. if (!_runLogFinalizer.IsHostStopping)
  71. await MarkRunFailedAsync(runLogId, now, ex.Message);
  72. throw;
  73. }
  74. finally
  75. {
  76. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  77. }
  78. }
  79. /// <summary>
  80. /// S7 成品入库双模式入站:执行器抽主/明细 → stg,再跑 WOI 标准层(读 stg)。
  81. /// </summary>
  82. public async Task<ProductionReceiptInboundResult> RunInboundAsync(
  83. long tenantId = 0,
  84. bool fullRefresh = false,
  85. CancellationToken cancellationToken = default)
  86. {
  87. cancellationToken.ThrowIfCancellationRequested();
  88. await EnsureTablesAsync();
  89. await EnsureStgTableAsync();
  90. var now = DateTime.Now;
  91. // 入参 tenantId:本库源含 tenant_id 时可传 0(Writer 优先取源行)。
  92. var pullCtx = new MdpPullContext
  93. {
  94. TenantId = tenantId,
  95. FullRefresh = fullRefresh,
  96. TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND",
  97. BatchId = $"S7_NBR_IN_{now:yyyyMMddHHmmss}"
  98. };
  99. var pull = await PopulateStgAsync(pullCtx, cancellationToken);
  100. var batchId = $"S7_NBR_STD_{now:yyyyMMddHHmmss}";
  101. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  102. var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  103. try
  104. {
  105. result.Rows = await TransformStandardAsync(batchId, now);
  106. await MarkRunSuccessAsync(runLogId, now, result);
  107. }
  108. catch (Exception ex)
  109. {
  110. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  111. if (!_runLogFinalizer.IsHostStopping)
  112. await MarkRunFailedAsync(runLogId, now, ex.Message);
  113. throw;
  114. }
  115. finally
  116. {
  117. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  118. }
  119. return new ProductionReceiptInboundResult
  120. {
  121. PullBatchId = pullCtx.BatchId,
  122. RowsPulled = pull.RowsPulled,
  123. RowsWrittenStg = pull.RowsWritten,
  124. NewCursor = pull.NewCursor,
  125. TransformBatchId = result.BatchId,
  126. StdRows = result.Rows
  127. };
  128. }
  129. /// <summary>
  130. /// 双源切换 FULL Replace(Phase 1):从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 stg,
  131. /// 成功后在单事务内以「仅当前源」结果 FULL 重建 mdp_std_production_receipt(按 tenant 精确隔离),消除旧源独有业务键残留。
  132. /// 用 PullAllByEntityCodeAsync 抽尽 + FullRefresh=true;transform 仅当前 source_system;std DELETE+INSERT 同事务,Pull 成功后才进入。
  133. /// SQLSERVER 实体默认 status=0(就位不启用);dopdemorq 六表当前空,Phase 1 不实际执行 destructive 切换。
  134. /// </summary>
  135. public async Task<ProductionReceiptInboundResult> RunSourceSwitchFullAsync(
  136. string sourceCode = SqlServerSourceCode,
  137. string masterEntityCode = SqlServerMasterEntityCode,
  138. string detailEntityCode = SqlServerDetailEntityCode,
  139. long tenantId = 0,
  140. CancellationToken cancellationToken = default)
  141. {
  142. cancellationToken.ThrowIfCancellationRequested();
  143. await EnsureTablesAsync();
  144. await EnsureStgTableAsync();
  145. // 165 SQL Server 源无 tenant_id;显式 tenantId≤0 时经账套映射解析。
  146. tenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  147. var now = DateTime.Now;
  148. // 1) PullAll 抽尽 + fullRefresh=true。任一 Pull 失败会抛异常,std 未动。
  149. var pullCtx = new MdpPullContext
  150. {
  151. TenantId = tenantId,
  152. FullRefresh = true,
  153. TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND",
  154. BatchId = $"S7_NBR_SW_{now:yyyyMMddHHmmss}"
  155. };
  156. var master = await _pullDispatcher.PullAllByEntityCodeAsync(masterEntityCode, pullCtx, cancellationToken);
  157. var detail = await _pullDispatcher.PullAllByEntityCodeAsync(detailEntityCode, pullCtx, cancellationToken);
  158. // 2) FULL Replace:事务内 DELETE 当前 tenant 的 std,再 INSERT 仅当前 source_system 的结果。
  159. var batchId = $"S7_NBR_SWSTD_{now:yyyyMMddHHmmss}";
  160. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  161. var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  162. try
  163. {
  164. result.Rows = await MdpStdFullReplace.ReplaceAsync(
  165. _db, "mdp_std_production_receipt", tenantId, extraWhere: null,
  166. insertScopedAsync: () => TransformStandardAsync(batchId, now, sourceCode),
  167. cancellationToken);
  168. await MarkRunSuccessAsync(runLogId, now, result);
  169. }
  170. catch (Exception ex)
  171. {
  172. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  173. if (!_runLogFinalizer.IsHostStopping)
  174. await MarkRunFailedAsync(runLogId, now, ex.Message);
  175. throw;
  176. }
  177. finally
  178. {
  179. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  180. }
  181. return new ProductionReceiptInboundResult
  182. {
  183. PullBatchId = pullCtx.BatchId,
  184. RowsPulled = master.RowsPulled + detail.RowsPulled,
  185. RowsWrittenStg = master.RowsWritten + detail.RowsWritten,
  186. NewCursor = detail.NewCursor ?? master.NewCursor,
  187. TransformBatchId = batchId,
  188. StdRows = result.Rows
  189. };
  190. }
  191. private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
  192. MdpPullContext pullCtx, CancellationToken cancellationToken)
  193. {
  194. var master = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, pullCtx, cancellationToken);
  195. var detail = await _pullDispatcher.PullByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
  196. return (
  197. master.RowsPulled + detail.RowsPulled,
  198. master.RowsWritten + detail.RowsWritten,
  199. detail.NewCursor ?? master.NewCursor);
  200. }
  201. private async Task EnsureStgTableAsync()
  202. {
  203. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_production_receipt"));
  204. }
  205. /// <summary>防御式建表(与 UpdateScripts/1.0.214.sql 同构,幂等)。</summary>
  206. private async Task EnsureTablesAsync()
  207. {
  208. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_production_receipt"));
  209. }
  210. /// <summary>
  211. /// 标准化:stg(NbrMaster/NbrDetail) + 维表(本库) → mdp_std_production_receipt。
  212. /// </summary>
  213. private async Task<int> TransformStandardAsync(string batchId, DateTime now, string? sourceSystem = null, long tenantId = 0)
  214. {
  215. var gateSource = sourceSystem ?? MdpSourceIdentity.Native;
  216. if (!await _neutralGate.AllowsAsync(tenantId, "PROD_DOC", gateSource))
  217. return 0;
  218. // 双源:sourceSystem 非空时仅统计/转换当前 source(切源 FULL Replace 用);为空时保持既有全源行为不变。
  219. var srcClause = sourceSystem == null ? "" : " AND n.source_system=@Src AND d.source_system=@Src";
  220. var countPars = new List<SugarParameter>();
  221. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  222. var rows = await _db.Ado.GetIntAsync(
  223. $"""
  224. SELECT COUNT(1)
  225. FROM mdp_stg_production_receipt n
  226. INNER JOIN mdp_stg_production_receipt d
  227. ON d.source_table='NbrDetail'
  228. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  229. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
  230. WHERE n.source_table='NbrMaster'
  231. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
  232. AND LOWER(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive'))) IN ('1','true')
  233. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL{srcClause}
  234. """, countPars);
  235. var insertSql =
  236. $"""
  237. INSERT INTO mdp_std_production_receipt
  238. (tenant_id, factory_id, source_system, domain, master_rec_id, detail_rec_id, nbr, line,
  239. receipt_date, status, status_desc, remark, prod_line, work_ord, erp_work_ord,
  240. department, department_desc, applicant_name, item_num, item_name, item_spec, um,
  241. location_to, location_to_desc, lot_serial, qty_rec, qty_to, location_from, location_from_desc,
  242. ord_nbr, source_biz_key, sync_batch_id, sync_time)
  243. SELECT
  244. NULLIF(n.tenant_id, 0), 1, {MdpSourceIdentity.Resolve("n.source_table", "n.source_system")},
  245. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')),
  246. CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.RecID')) AS SIGNED),
  247. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS SIGNED),
  248. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')),
  249. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS SIGNED),
  250. STR_TO_DATE(REPLACE(LEFT(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Date')),'null'),''), 19), 'T', ' '), '%Y-%m-%d %H:%i:%s'),
  251. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  252. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  253. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Remark')),
  254. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ProdLine')),
  255. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.WorkOrd')),
  256. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Address')),
  257. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),
  258. TRIM(CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),''), ' ', IFNULL(p.Descr, ''))),
  259. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Name')),
  260. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')), i.Descr, i.Descr1, i.UM,
  261. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')), lt.descr,
  262. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LotSerial')) AS CHAR),
  263. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyRec')) AS DECIMAL(18,5)),
  264. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyTo')) AS DECIMAL(18,5)),
  265. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')), lf.descr,
  266. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')),
  267. CONCAT(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), '#',
  268. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID'))),
  269. @BatchId, @Now
  270. FROM mdp_stg_production_receipt n
  271. INNER JOIN mdp_stg_production_receipt d
  272. ON d.source_table='NbrDetail'
  273. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  274. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
  275. LEFT JOIN DepartmentMaster p
  276. ON p.Domain = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain'))
  277. AND p.Department = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department'))
  278. LEFT JOIN ItemMaster i
  279. ON i.Domain = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  280. AND i.ItemNum = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum'))
  281. LEFT JOIN LocationMaster lt
  282. ON lt.tenant_id = IFNULL(d.tenant_id, 0)
  283. AND lt.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo'))
  284. LEFT JOIN LocationMaster lf
  285. ON lf.tenant_id = IFNULL(d.tenant_id, 0)
  286. AND lf.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom'))
  287. WHERE n.source_table='NbrMaster'
  288. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
  289. AND LOWER(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive'))) IN ('1','true')
  290. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL{srcClause}
  291. AND NULLIF(n.tenant_id, 0) IS NOT NULL
  292. ON DUPLICATE KEY UPDATE
  293. factory_id=VALUES(factory_id), master_rec_id=VALUES(master_rec_id), nbr=VALUES(nbr), line=VALUES(line),
  294. receipt_date=VALUES(receipt_date), status=VALUES(status), status_desc=VALUES(status_desc), remark=VALUES(remark),
  295. prod_line=VALUES(prod_line), work_ord=VALUES(work_ord), erp_work_ord=VALUES(erp_work_ord),
  296. department=VALUES(department), department_desc=VALUES(department_desc), applicant_name=VALUES(applicant_name),
  297. item_num=VALUES(item_num), item_name=VALUES(item_name), item_spec=VALUES(item_spec), um=VALUES(um),
  298. location_to=VALUES(location_to), location_to_desc=VALUES(location_to_desc), lot_serial=VALUES(lot_serial),
  299. qty_rec=VALUES(qty_rec), qty_to=VALUES(qty_to), location_from=VALUES(location_from),
  300. location_from_desc=VALUES(location_from_desc), ord_nbr=VALUES(ord_nbr),
  301. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  302. """;
  303. var insPars = new List<SugarParameter>
  304. {
  305. new("@BatchId", batchId),
  306. new("@Now", now)
  307. };
  308. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  309. await MdpSchemaAligner.ExecuteAsync(_db, insertSql, insPars);
  310. return rows;
  311. }
  312. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  313. {
  314. await MdpSchemaAligner.ExecuteAsync(_db,
  315. """
  316. INSERT INTO mdp_transform_run_log
  317. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  318. VALUES (0, @JobCode, 'S5生产入库单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  319. """,
  320. new SugarParameter("@JobCode", JobCode),
  321. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  322. new SugarParameter("@BatchId", batchId),
  323. new SugarParameter("@StartTime", startedAt));
  324. return await _db.Ado.GetLongAsync(
  325. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  326. new List<SugarParameter> { new("@BatchId", batchId) });
  327. }
  328. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, ProductionReceiptMdpSyncResult result)
  329. {
  330. var finishedAt = DateTime.Now;
  331. await MdpSchemaAligner.ExecuteAsync(_db,
  332. """
  333. UPDATE mdp_transform_run_log
  334. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  335. stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  336. WHERE id=@Id
  337. """,
  338. new SugarParameter("@EndTime", finishedAt),
  339. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  340. new SugarParameter("@StandardRows", result.Rows),
  341. new SugarParameter("@Id", runLogId));
  342. }
  343. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  344. {
  345. var finishedAt = DateTime.Now;
  346. await MdpSchemaAligner.ExecuteAsync(_db,
  347. """
  348. UPDATE mdp_transform_run_log
  349. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  350. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  351. WHERE id=@Id
  352. """,
  353. new SugarParameter("@EndTime", finishedAt),
  354. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  355. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  356. new SugarParameter("@Id", runLogId));
  357. }
  358. private static string NormalizeTriggerType(string? triggerType)
  359. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  360. }
  361. /// <summary>生产入库单 MDP 同步转换结果。</summary>
  362. public sealed class ProductionReceiptMdpSyncResult
  363. {
  364. public long RunLogId { get; set; }
  365. public string BatchId { get; set; } = string.Empty;
  366. public int Rows { get; set; }
  367. }
  368. /// <summary>S7 成品入库双模式入站结果。</summary>
  369. public sealed class ProductionReceiptInboundResult
  370. {
  371. public string PullBatchId { get; set; } = string.Empty;
  372. public int RowsPulled { get; set; }
  373. public int RowsWrittenStg { get; set; }
  374. public string? NewCursor { get; set; }
  375. public string TransformBatchId { get; set; } = string.Empty;
  376. public int StdRows { get; set; }
  377. }