ProductionIssueMdpSyncService.cs 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474
  1. using Admin.NET.Plugin.AiDOP.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  3. using Admin.NET.Plugin.AiDOP.Infrastructure;
  4. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  5. namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  6. /// <summary>
  7. /// S5 生产领料单 数据中台只读同步转换服务(DOP 内部,头级单表)。
  8. ///
  9. /// B-① 双源通用管线(2026-07-29 由 bespoke 直读 NbrMaster 迁入):
  10. /// 源(NbrMaster) --MdpDbPullExecutor--> mdp_stg_production_issue(raw_data JSON) --transform--> mdp_std_production_issue。
  11. /// 入站实体:S5_PRODUCTION_ISSUE_MASTER(源=AIDOPDEV_MYSQL) / S5_PRODUCTION_ISSUE_MASTER_SQLSERVER(源=DOPDEMORQ_SQLSERVER)。
  12. /// 维表 DepartmentMaster 仍本库 LEFT JOIN(与 mode-A 采购/生产收货同构;department_desc 富化,部门 code 由事实表保留)。
  13. ///
  14. /// 约束:
  15. /// - 只读源/贴源,仅写 mdp_stg_production_issue / mdp_std_production_issue;绝不写 NbrMaster;不碰 WorkOrderPickBillService 写路径。
  16. /// - 头过滤 Type='SM' AND IsActive(领料无 IsReturn 语义);SM=领料单业务类型(Type SM/WOI/WOD/CA 互斥)。
  17. /// - 跨源类型兼容经 MdpJsonSql(datetime ISO-T / bit true-false / JSON null 归一),不改 MDP 核心。
  18. /// - 切源单 active source + FULL Replace(决策1A);不做 merge、不做回写/outbox(决策2A)。
  19. /// - SM 为 0 行时转换成功完成、处理数为 0,不报错。
  20. /// </summary>
  21. public class ProductionIssueMdpSyncService : ITransient
  22. {
  23. private const string JobCode = "S5_PRODUCTION_ISSUE_MDP_SYNC";
  24. private const string InboundEntityCode = "S5_PRODUCTION_ISSUE_MASTER";
  25. // 双源(dopdemorq SQL Server):第二源与其入站实体。实体默认 status=0(就位不启用)。
  26. private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
  27. private const string SqlServerEntityCode = "S5_PRODUCTION_ISSUE_MASTER_SQLSERVER";
  28. /// <summary>165 SQL Server 源无 tenant_id 列;切源时须经账套映射解析。</summary>
  29. private const string SourceZtid = "pbxfxp";
  30. private readonly ISqlSugarClient _db;
  31. private readonly TransformRunLogFinalizer _runLogFinalizer;
  32. private readonly MdpSourcePullDispatcher _pullDispatcher;
  33. private readonly MdpNeutralSourceGate _neutralGate;
  34. public ProductionIssueMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher, TransformRunLogFinalizer runLogFinalizer, MdpNeutralSourceGate neutralGate)
  35. {
  36. _db = db;
  37. _runLogFinalizer = runLogFinalizer;
  38. _pullDispatcher = pullDispatcher;
  39. _neutralGate = neutralGate;
  40. }
  41. /// <summary>全量:按当前启用入站实体灌 stg → 标准层(WP9 S5 后默认 165)。</summary>
  42. public async Task<ProductionIssueMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  43. {
  44. cancellationToken.ThrowIfCancellationRequested();
  45. await EnsureTablesAsync();
  46. await EnsureStgTableAsync();
  47. var active = await ResolveActiveInboundAsync(cancellationToken);
  48. var now = DateTime.Now;
  49. var batchId = $"S5_PROD_ISSUE_FULL_{now:yyyyMMddHHmmss}";
  50. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  51. var result = new ProductionIssueMdpSyncResult { 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 = "S5_PRODUCTION_ISSUE_INBOUND",
  60. BatchId = $"{batchId}_PULL"
  61. };
  62. await PopulateStgAsync(active.EntityCode, pullCtx, cancellationToken);
  63. if (active.UseSourceFilter)
  64. {
  65. result.HeadRows = await MdpStdFullReplace.ReplaceAsync(
  66. _db, "mdp_std_production_issue", 0, extraWhere: null,
  67. insertScopedAsync: () => TransformHeadStandardAsync(batchId, now, active.SourceCode),
  68. cancellationToken);
  69. }
  70. else
  71. {
  72. result.HeadRows = await TransformHeadStandardAsync(batchId, now);
  73. }
  74. await MarkRunSuccessAsync(runLogId, now, result);
  75. return result;
  76. }
  77. catch (Exception ex)
  78. {
  79. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  80. if (!_runLogFinalizer.IsHostStopping)
  81. await MarkRunFailedAsync(runLogId, now, ex.Message);
  82. throw;
  83. }
  84. finally
  85. {
  86. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  87. }
  88. }
  89. /// <summary>双模式入站:执行器抽头落 stg,再跑标准层(读 stg)。</summary>
  90. public async Task<ProductionIssueInboundResult> RunInboundAsync(
  91. long tenantId = 0, bool fullRefresh = false, CancellationToken cancellationToken = default)
  92. {
  93. cancellationToken.ThrowIfCancellationRequested();
  94. await EnsureTablesAsync();
  95. await EnsureStgTableAsync();
  96. var active = await ResolveActiveInboundAsync(cancellationToken);
  97. var now = DateTime.Now;
  98. // 165 源无 tenant_id 须映射;本库 NbrMaster 含 tenant_id,tenantId=0 时 Writer 优先取源行。
  99. var ctxTenantId = active.UseSourceFilter
  100. ? AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId)
  101. : tenantId;
  102. var pullCtx = new MdpPullContext
  103. {
  104. TenantId = ctxTenantId,
  105. FullRefresh = fullRefresh,
  106. TaskCode = "S5_PRODUCTION_ISSUE_INBOUND",
  107. BatchId = $"S5_PROD_ISSUE_IN_{now:yyyyMMddHHmmss}"
  108. };
  109. var pull = await PopulateStgAsync(active.EntityCode, pullCtx, cancellationToken);
  110. var batchId = $"S5_PROD_ISSUE_STD_{now:yyyyMMddHHmmss}";
  111. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  112. var result = new ProductionIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  113. try
  114. {
  115. if (active.UseSourceFilter)
  116. {
  117. result.HeadRows = await MdpStdFullReplace.ReplaceAsync(
  118. _db, "mdp_std_production_issue", ctxTenantId, extraWhere: null,
  119. insertScopedAsync: () => TransformHeadStandardAsync(batchId, now, active.SourceCode),
  120. cancellationToken);
  121. }
  122. else
  123. {
  124. result.HeadRows = await TransformHeadStandardAsync(batchId, now);
  125. }
  126. await MarkRunSuccessAsync(runLogId, now, result);
  127. }
  128. catch (Exception ex)
  129. {
  130. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  131. if (!_runLogFinalizer.IsHostStopping)
  132. await MarkRunFailedAsync(runLogId, now, ex.Message);
  133. throw;
  134. }
  135. finally
  136. {
  137. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  138. }
  139. return new ProductionIssueInboundResult
  140. {
  141. PullBatchId = pullCtx.BatchId,
  142. RowsPulled = pull.RowsPulled,
  143. RowsWrittenStg = pull.RowsWritten,
  144. NewCursor = pull.NewCursor,
  145. TransformBatchId = result.BatchId,
  146. StdRows = result.HeadRows
  147. };
  148. }
  149. /// <summary>
  150. /// WP9 S5:切到 165 权威源并 FULL Replace 标准层;同时确保本库历史 SM 已打 AIDOP_LEGACY。
  151. /// </summary>
  152. public async Task<ProductionIssueInboundResult> ActivateSqlServerSourceAndSwitchAsync(
  153. long tenantId = 0,
  154. CancellationToken cancellationToken = default)
  155. {
  156. cancellationToken.ThrowIfCancellationRequested();
  157. await EnsureNbrMasterSourceSystemColumnAsync();
  158. await MarkLocalSmLegacyAsync();
  159. await FlipInboundEntityStatusAsync();
  160. return await RunSourceSwitchFullAsync(
  161. SqlServerSourceCode, SqlServerEntityCode, tenantId, cancellationToken);
  162. }
  163. private async Task<(string EntityCode, string SourceCode, bool UseSourceFilter)> ResolveActiveInboundAsync(
  164. CancellationToken cancellationToken)
  165. {
  166. var sqlOn = await _db.Queryable<MdpEntity>()
  167. .Where(x => x.EntityCode == SqlServerEntityCode && x.Status == 1)
  168. .AnyAsync(cancellationToken);
  169. if (sqlOn)
  170. return (SqlServerEntityCode, SqlServerSourceCode, true);
  171. var mysqlOn = await _db.Queryable<MdpEntity>()
  172. .Where(x => x.EntityCode == InboundEntityCode && x.Status == 1)
  173. .AnyAsync(cancellationToken);
  174. if (mysqlOn)
  175. return (InboundEntityCode, "AIDOPDEV_MYSQL", false);
  176. throw Oops.Oh("生产领料同步失败:未启用任何 S5 入站实体,请执行 WP9 S5 切源或检查 mdp_entity。");
  177. }
  178. private async Task FlipInboundEntityStatusAsync()
  179. {
  180. await MdpSchemaAligner.ExecuteAsync(_db,
  181. """
  182. UPDATE mdp_entity
  183. SET status = 0, update_time = NOW()
  184. WHERE entity_code = @Mysql AND status <> 0
  185. """,
  186. new SugarParameter("@Mysql", InboundEntityCode));
  187. await MdpSchemaAligner.ExecuteAsync(_db,
  188. """
  189. UPDATE mdp_entity
  190. SET status = 1, update_time = NOW()
  191. WHERE entity_code = @Sql AND status <> 1
  192. """,
  193. new SugarParameter("@Sql", SqlServerEntityCode));
  194. }
  195. private async Task EnsureNbrMasterSourceSystemColumnAsync()
  196. {
  197. var exists = await _db.Ado.GetIntAsync(
  198. """
  199. SELECT COUNT(1) FROM information_schema.COLUMNS
  200. WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'NbrMaster' AND COLUMN_NAME = 'source_system'
  201. """);
  202. if (exists == 0)
  203. {
  204. await MdpSchemaAligner.ExecuteAsync(_db,
  205. "ALTER TABLE NbrMaster ADD COLUMN source_system VARCHAR(50) NULL DEFAULT NULL AFTER Remark");
  206. }
  207. }
  208. private async Task MarkLocalSmLegacyAsync()
  209. {
  210. await MdpSchemaAligner.ExecuteAsync(_db,
  211. """
  212. UPDATE NbrMaster
  213. SET source_system = 'AIDOP_NATIVE'
  214. WHERE Type = 'SM'
  215. AND IFNULL(source_system, '') IN ('', 'AIDOP', 'AIDOP_LEGACY')
  216. """);
  217. }
  218. /// <summary>
  219. /// 双源切换 FULL Replace(Phase 1):从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 stg,
  220. /// 成功后在单事务内以「仅当前源」结果 FULL 重建 mdp_std_production_issue(按 tenant 精确隔离),消除旧源独有业务键残留。
  221. /// PullAllByEntityCodeAsync 抽尽 + FullRefresh=true;transform 仅当前 source_system;Pull 成功后才进入 destructive replace。
  222. /// SQLSERVER 实体默认 status=0(就位不启用);本批 FULL-only,不启自动增量。
  223. /// </summary>
  224. public async Task<ProductionIssueInboundResult> RunSourceSwitchFullAsync(
  225. string sourceCode = SqlServerSourceCode,
  226. string entityCode = SqlServerEntityCode,
  227. long tenantId = 0,
  228. CancellationToken cancellationToken = default)
  229. {
  230. cancellationToken.ThrowIfCancellationRequested();
  231. await EnsureTablesAsync();
  232. await EnsureStgTableAsync();
  233. // 165 SQL Server 源无 tenant_id;显式 tenantId≤0 时经账套映射解析。
  234. tenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  235. var now = DateTime.Now;
  236. var pullCtx = new MdpPullContext
  237. {
  238. TenantId = tenantId,
  239. FullRefresh = true,
  240. TaskCode = "S5_PRODUCTION_ISSUE_INBOUND",
  241. BatchId = $"S5_PROD_ISSUE_SW_{now:yyyyMMddHHmmss}"
  242. };
  243. var pull = await _pullDispatcher.PullAllByEntityCodeAsync(entityCode, pullCtx, cancellationToken);
  244. var batchId = $"S5_PROD_ISSUE_SWSTD_{now:yyyyMMddHHmmss}";
  245. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  246. var result = new ProductionIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  247. try
  248. {
  249. result.HeadRows = await MdpStdFullReplace.ReplaceAsync(
  250. _db, "mdp_std_production_issue", tenantId, extraWhere: null,
  251. insertScopedAsync: () => TransformHeadStandardAsync(batchId, now, sourceCode),
  252. cancellationToken);
  253. await MarkRunSuccessAsync(runLogId, now, result);
  254. }
  255. catch (Exception ex)
  256. {
  257. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  258. if (!_runLogFinalizer.IsHostStopping)
  259. await MarkRunFailedAsync(runLogId, now, ex.Message);
  260. throw;
  261. }
  262. finally
  263. {
  264. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  265. }
  266. return new ProductionIssueInboundResult
  267. {
  268. PullBatchId = pullCtx.BatchId,
  269. RowsPulled = pull.RowsPulled,
  270. RowsWrittenStg = pull.RowsWritten,
  271. NewCursor = pull.NewCursor,
  272. TransformBatchId = batchId,
  273. StdRows = result.HeadRows
  274. };
  275. }
  276. private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
  277. string entityCode, MdpPullContext pullCtx, CancellationToken cancellationToken)
  278. {
  279. var r = await _pullDispatcher.PullByEntityCodeAsync(entityCode, pullCtx, cancellationToken);
  280. return (r.RowsPulled, r.RowsWritten, r.NewCursor);
  281. }
  282. private async Task EnsureStgTableAsync()
  283. {
  284. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_production_issue"));
  285. }
  286. /// <summary>防御式建表(与 UpdateScripts DDL 同构,幂等)。</summary>
  287. private async Task EnsureTablesAsync()
  288. {
  289. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_production_issue"));
  290. }
  291. /// <summary>
  292. /// 标准化头:stg(NbrMaster,SM) + DepartmentMaster(本库) → mdp_std_production_issue。
  293. /// sourceSystem 非空时仅转当前源(切源 FULL Replace 用);为空时全源(本库全量)。跨源类型经 MdpJsonSql 兼容。
  294. /// </summary>
  295. private async Task<int> TransformHeadStandardAsync(string batchId, DateTime now, string? sourceSystem = null, long tenantId = 0)
  296. {
  297. var gateSource = sourceSystem ?? MdpSourceIdentity.Native;
  298. if (!await _neutralGate.AllowsAsync(tenantId, "PROD_DOC", gateSource))
  299. return 0;
  300. var srcClause = sourceSystem == null ? "" : " AND m.source_system=@Src";
  301. // 表达式片段(跨源类型兼容)
  302. var mTenant = MdpJsonSql.TenantFromStgJsonCol("m", MdpJsonSql.Int("m", "tenant_id"));
  303. var domainE = MdpJsonSql.Str("m", "Domain");
  304. var deptE = MdpJsonSql.Str("m", "Department");
  305. var statusUpper = $"UPPER(IFNULL({MdpJsonSql.Str("m", "Status")},''))";
  306. var pretreatE = MdpJsonSql.Str("m", "PretreatmentState");
  307. var transTypeE = MdpJsonSql.Str("m", "TransType");
  308. var countPars = new List<SugarParameter>();
  309. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  310. var rows = await _db.Ado.GetIntAsync(
  311. $"SELECT COUNT(1) FROM mdp_stg_production_issue m WHERE m.source_table='NbrMaster' " +
  312. $"AND {MdpJsonSql.Str("m", "Type")}='SM' AND {MdpJsonSql.BoolTrue("m", "IsActive")}{srcClause}",
  313. countPars);
  314. var insertSql =
  315. $"""
  316. INSERT INTO mdp_std_production_issue
  317. (tenant_id, factory_id, source_system, domain, rec_id, nbr, issue_date, status, status_desc,
  318. work_ord, department, department_desc, qty_ord, prod_line, applicant_name, issue_user,
  319. user1, remark, create_user, source_create_time, pretreatment_state, trans_type, trans_type_text,
  320. eff_date, address, source_biz_key, sync_batch_id, sync_time)
  321. SELECT
  322. {mTenant}, 1, {MdpSourceIdentity.Resolve("m.source_table", "m.source_system")},
  323. {domainE}, {MdpJsonSql.Int("m", "RecID")}, {MdpJsonSql.Str("m", "Nbr")},
  324. {MdpJsonSql.DateTimeSec("m", "Date")},
  325. {statusUpper},
  326. CASE
  327. WHEN IFNULL({pretreatE},'')<>'' THEN {pretreatE}
  328. WHEN {statusUpper}='C' THEN '已下架'
  329. WHEN {statusUpper}='A' THEN '备料中'
  330. ELSE ''
  331. END,
  332. {MdpJsonSql.Str("m", "WorkOrd")}, {deptE},
  333. TRIM(CONCAT(IFNULL({deptE},''), ' ', IFNULL(d.Descr, ''))),
  334. {MdpJsonSql.Dec("m", "QtyOrd", 18, 5)}, {MdpJsonSql.Str("m", "ProdLine")}, {MdpJsonSql.Str("m", "Name")},
  335. {MdpJsonSql.Str("m", "User1")},
  336. {MdpJsonSql.Str("m", "User1")}, {MdpJsonSql.Str("m", "Remark")}, {MdpJsonSql.Str("m", "CreateUser")},
  337. {MdpJsonSql.DateTimeSec("m", "CreateTime")},
  338. IFNULL({pretreatE}, ''),
  339. CASE WHEN {transTypeE}='Z61' THEN '补料' ELSE '正常' END,
  340. CASE WHEN {transTypeE}='PrevProcess' THEN '需要前处理' ELSE '' END,
  341. {MdpJsonSql.DateTimeSec("m", "EffDate")}, {MdpJsonSql.Str("m", "Address")},
  342. IFNULL(NULLIF(m.source_biz_key,''), CONCAT(IFNULL({domainE},''), ':', IFNULL({MdpJsonSql.Str("m", "RecID")},''))),
  343. @BatchId, @Now
  344. FROM mdp_stg_production_issue m
  345. LEFT JOIN DepartmentMaster d ON d.Domain = {domainE} AND d.Department = {deptE}
  346. WHERE m.source_table='NbrMaster'
  347. AND {MdpJsonSql.Str("m", "Type")}='SM'
  348. AND {MdpJsonSql.BoolTrue("m", "IsActive")}{srcClause}
  349. AND {MdpJsonSql.TenantGuard(mTenant)}
  350. ON DUPLICATE KEY UPDATE
  351. factory_id=VALUES(factory_id), source_system=VALUES(source_system), nbr=VALUES(nbr), issue_date=VALUES(issue_date),
  352. status=VALUES(status), status_desc=VALUES(status_desc), work_ord=VALUES(work_ord),
  353. department=VALUES(department), department_desc=VALUES(department_desc), qty_ord=VALUES(qty_ord),
  354. prod_line=VALUES(prod_line), applicant_name=VALUES(applicant_name), issue_user=VALUES(issue_user),
  355. user1=VALUES(user1), remark=VALUES(remark), create_user=VALUES(create_user),
  356. source_create_time=VALUES(source_create_time), pretreatment_state=VALUES(pretreatment_state),
  357. trans_type=VALUES(trans_type), trans_type_text=VALUES(trans_type_text),
  358. eff_date=VALUES(eff_date), address=VALUES(address),
  359. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  360. """;
  361. var insPars = new List<SugarParameter>
  362. {
  363. new("@BatchId", batchId),
  364. new("@Now", now)
  365. };
  366. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  367. await MdpSchemaAligner.ExecuteAsync(_db, insertSql, insPars);
  368. return rows;
  369. }
  370. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  371. {
  372. await MdpSchemaAligner.ExecuteAsync(_db,
  373. """
  374. INSERT INTO mdp_transform_run_log
  375. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  376. VALUES (0, @JobCode, 'S5生产领料单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  377. """,
  378. new SugarParameter("@JobCode", JobCode),
  379. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  380. new SugarParameter("@BatchId", batchId),
  381. new SugarParameter("@StartTime", startedAt));
  382. return await _db.Ado.GetLongAsync(
  383. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  384. new List<SugarParameter> { new("@BatchId", batchId) });
  385. }
  386. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, ProductionIssueMdpSyncResult result)
  387. {
  388. var finishedAt = DateTime.Now;
  389. await MdpSchemaAligner.ExecuteAsync(_db,
  390. """
  391. UPDATE mdp_transform_run_log
  392. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  393. stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  394. WHERE id=@Id
  395. """,
  396. new SugarParameter("@EndTime", finishedAt),
  397. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  398. new SugarParameter("@StandardRows", result.HeadRows),
  399. new SugarParameter("@Id", runLogId));
  400. }
  401. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  402. {
  403. var finishedAt = DateTime.Now;
  404. await MdpSchemaAligner.ExecuteAsync(_db,
  405. """
  406. UPDATE mdp_transform_run_log
  407. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  408. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  409. WHERE id=@Id
  410. """,
  411. new SugarParameter("@EndTime", finishedAt),
  412. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  413. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  414. new SugarParameter("@Id", runLogId));
  415. }
  416. private static string NormalizeTriggerType(string? triggerType)
  417. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  418. }
  419. /// <summary>生产领料单 MDP 同步转换结果。</summary>
  420. public sealed class ProductionIssueMdpSyncResult
  421. {
  422. public long RunLogId { get; set; }
  423. public string BatchId { get; set; } = string.Empty;
  424. public int HeadRows { get; set; }
  425. }
  426. /// <summary>生产领料双源入站结果(stg 抽数 + std 转换)。</summary>
  427. public sealed class ProductionIssueInboundResult
  428. {
  429. public string PullBatchId { get; set; } = string.Empty;
  430. public int RowsPulled { get; set; }
  431. public int RowsWrittenStg { get; set; }
  432. public string? NewCursor { get; set; }
  433. public string TransformBatchId { get; set; } = string.Empty;
  434. public int StdRows { get; set; }
  435. }