ProductionReceiptMdpSyncService.cs 22 KB

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