PurchaseOrderCompletionMdpSyncService.cs 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690
  1. using Microsoft.Extensions.Hosting;
  2. using Microsoft.Extensions.Logging;
  3. using Microsoft.Extensions.Options;
  4. using SqlSugar;
  5. using System.Data;
  6. using System.Text;
  7. using Admin.NET.Plugin.AiDOP.DataPlatform;
  8. namespace Admin.NET.Plugin.AiDOP.Supply;
  9. /// <summary>
  10. /// S8 Stage-3 采购订单完成态标准事实层投影服务。
  11. ///
  12. /// <para>链路:Source B(WMS/MES 的 <c>PurOrdDetail</c>)→ <c>mdp_std_purchase_order_completion</c>。
  13. /// 目的是让 S8 Stage-3 能在【不覆盖】<c>mdp_std_purchase_order.status</c> 的前提下拿到
  14. /// 「这条采购行收货是否已关闭」这一 WMS 侧事实。</para>
  15. ///
  16. /// <para><b>Grain = 一条采购行</b>(<c>tenant_id + po_no + po_line</c>)。
  17. /// <c>domain</c> / <c>potype</c> / <c>source_row_id</c> / <c>source_id</c> 是 provenance,不进唯一键。</para>
  18. ///
  19. /// <para><b>当前状态,不是事件</b>:上游收货过程的写法是
  20. /// <c>Status = (case when 收满 then 'C' else '' end)</c>,退货会把 <c>'C'</c> 打回空串。
  21. /// 因此本服务的 UPSERT 必须无条件刷新结论,允许 COMPLETED → NOT_COMPLETED 回退。</para>
  22. ///
  23. /// <para><b>Full Reconciliation 是正确性主路径</b>:HotWatch 只覆盖活动中的采购单,
  24. /// 实测 5 个 <c>PUR_ORDER</c> 关注里 3 个已提前终止(且无一是因采购完成而终止),
  25. /// 98 条源行里 88 条根本没有对应关注。所以单靠增量必然漏判,周期性全量对账不可省。</para>
  26. ///
  27. /// <para><b>刻意不做</b>:不读 IQC、不算 ObjectCompleted、不做订单级聚合、
  28. /// 不用数量重新推导完成(唯一 Authority 是源侧 Status)、不写任何源表、不改 Source A 的任何列。</para>
  29. /// </summary>
  30. public class PurchaseOrderCompletionMdpSyncService : ITransient
  31. {
  32. private const string JobCode = "S3_PO_COMPLETION_MDP_SYNC";
  33. private const string JobName = "S3采购订单完成态标准事实层投影";
  34. /// <summary>完成结论取值域。</summary>
  35. public const string StatusCompleted = "COMPLETED";
  36. public const string StatusNotCompleted = "NOT_COMPLETED";
  37. public const string StatusUnknown = "UNKNOWN";
  38. private readonly ISqlSugarClient _db;
  39. private readonly TransformRunLogFinalizer _runLogFinalizer;
  40. private readonly DataPlatform.MdpSourceScopeFactory _scopeFactory;
  41. private readonly AidopPoCompletionOptions _opt;
  42. private readonly IHostEnvironment _env;
  43. private readonly ILogger<PurchaseOrderCompletionMdpSyncService> _logger;
  44. private readonly MdpNeutralSourceGate _neutralGate;
  45. public PurchaseOrderCompletionMdpSyncService(
  46. ISqlSugarClient db,
  47. DataPlatform.MdpSourceScopeFactory scopeFactory,
  48. IOptions<AidopPoCompletionOptions> opt,
  49. IHostEnvironment env,
  50. ILogger<PurchaseOrderCompletionMdpSyncService> logger,
  51. TransformRunLogFinalizer runLogFinalizer,
  52. MdpNeutralSourceGate neutralGate)
  53. {
  54. _db = db;
  55. _runLogFinalizer = runLogFinalizer;
  56. _scopeFactory = scopeFactory;
  57. _opt = opt.Value;
  58. _env = env;
  59. _logger = logger;
  60. _neutralGate = neutralGate;
  61. }
  62. /// <summary>
  63. /// 完成态映射。<b>这是全链路唯一的状态映射入口</b>——HotWatch 增量若接入,必须复用本方法,
  64. /// 不得另写一套。
  65. ///
  66. /// <para><c>'C'</c> → COMPLETED;空串 → NOT_COMPLETED;
  67. /// <c>null</c> 与任何其它非空值 → UNKNOWN。</para>
  68. ///
  69. /// <para>注意这里刻意<b>不</b>写成 <c>IFNULL(status,'') == ""</c>:
  70. /// NULL 表示「源侧没给值」,未知值表示「源侧给了我们不认识的值」,
  71. /// 两者都不等于「已确认尚未完成」,滑成 NOT_COMPLETED 会让下游把未知当成确定结论。</para>
  72. /// </summary>
  73. public static string MapCompletionStatus(string? rawStatus)
  74. {
  75. if (rawStatus is null) return StatusUnknown;
  76. var s = rawStatus.Trim();
  77. if (s.Length == 0) return StatusNotCompleted;
  78. return string.Equals(s, "C", StringComparison.OrdinalIgnoreCase)
  79. ? StatusCompleted
  80. : StatusUnknown;
  81. }
  82. /// <summary>写入本表时使用的 raw_status 归一:保留 NULL,其余 TRIM。</summary>
  83. public static string? NormalizeRawStatus(string? rawStatus) => rawStatus?.Trim();
  84. /// <summary>
  85. /// 全量对账(Full Reconciliation)。
  86. /// <paramref name="tenantId"/> 为 0 表示全租户;&gt;0 时只写/只淘汰该租户,绝不触碰其它租户。
  87. /// </summary>
  88. public async Task<PurchaseOrderCompletionSyncResult> RunFullAsync(
  89. long tenantId = 0,
  90. string triggerType = "AUTO",
  91. CancellationToken cancellationToken = default)
  92. {
  93. cancellationToken.ThrowIfCancellationRequested();
  94. var now = DateTime.Now;
  95. var batchId = $"S3_PO_COMPL_{(tenantId > 0 ? tenantId + "_" : "")}{now:yyyyMMddHHmmss}";
  96. var runLogId = await InsertRunLogAsync(batchId, now, triggerType, tenantId);
  97. var result = new PurchaseOrderCompletionSyncResult { BatchId = batchId, RunLogId = runLogId };
  98. try
  99. {
  100. // ① 源绑定:非开发环境必须显式绑定真实 WMS 源,绝不静默回落
  101. var binding = await ResolveSourceBindingAsync(cancellationToken);
  102. result.SourceCode = binding.SourceCode;
  103. var remote = await _scopeFactory.GetScopeAsync(binding.SourceCode, cancellationToken);
  104. // 作用域工厂对「本库样板源」与「未配账号的源」会直接返回主库连接。
  105. // 那会让本服务把 Source A 的本地表当成 WMS 执行态读进来,结论会整体错误,
  106. // 所以这里必须结构性挡住,而不是依赖配置写对。
  107. if (ReferenceEquals(remote, _db))
  108. throw new InvalidOperationException(
  109. $"源 {binding.SourceCode} 解析结果是主库连接,说明它是本库样板源或未配置独立账号,"
  110. + "不能承担 S8_PO_COMPLETION_SOURCE 角色");
  111. // ② 读源:只取本角色范围内的采购行
  112. var sourceRows = await ReadSourceRowsAsync(remote, cancellationToken);
  113. result.SourceRows = sourceRows.Count;
  114. // ③ 归属与校验
  115. var tenantIndex = await LoadLocalTenantIndexAsync(cancellationToken);
  116. var poLineIndex = await LoadLocalPoLineIndexAsync(cancellationToken);
  117. var pending = new List<CompletionRow>(sourceRows.Count);
  118. foreach (var row in sourceRows)
  119. {
  120. cancellationToken.ThrowIfCancellationRequested();
  121. if (string.IsNullOrWhiteSpace(row.PoNo) || string.IsNullOrWhiteSpace(row.PoLine))
  122. {
  123. result.SkippedRows++;
  124. continue;
  125. }
  126. var resolved = ResolveTenant(tenantIndex, row.PoNo, row.Domain);
  127. if (resolved <= 0)
  128. {
  129. result.SkippedRows++;
  130. result.TenantUnresolvedRows++;
  131. _logger.LogWarning(
  132. "[PoCompletion] 跳过:采购单 {PurOrd} 行 {Line} 无法唯一解析租户(候选数 {Count})",
  133. row.PoNo, row.PoLine, CountTenantCandidates(tenantIndex, row.PoNo, row.Domain));
  134. continue;
  135. }
  136. // 单租户刷新:其它租户的源行本轮不参与,也因此不能被本轮淘汰
  137. if (tenantId > 0 && resolved != tenantId)
  138. continue;
  139. var matches = poLineIndex.TryGetValue(LineKey(resolved, row.PoNo, row.PoLine), out var c) ? c : 0;
  140. if (matches > 1)
  141. throw new InvalidOperationException(
  142. $"数据契约破坏:租户 {resolved} 的采购行 {row.PoNo}#{row.PoLine} 在本库命中 {matches} 条,"
  143. + "唯一键 uk_po_line 应当保证唯一");
  144. if (matches == 0)
  145. {
  146. result.SkippedRows++;
  147. result.UnmatchedRows++;
  148. _logger.LogWarning(
  149. "[PoCompletion] 跳过:源行 {PurOrd}#{Line} 在本库租户 {Tenant} 下无对应采购行",
  150. row.PoNo, row.PoLine, resolved);
  151. continue;
  152. }
  153. result.MatchedRows++;
  154. var raw = NormalizeRawStatus(row.RawStatus);
  155. var completion = MapCompletionStatus(row.RawStatus);
  156. if (completion == StatusCompleted) result.CompletedRows++;
  157. else if (completion == StatusNotCompleted) result.NotCompletedRows++;
  158. else result.UnknownStatusRows++;
  159. pending.Add(new CompletionRow
  160. {
  161. TenantId = resolved,
  162. PoNo = row.PoNo.Trim(),
  163. PoLine = row.PoLine.Trim(),
  164. Domain = string.IsNullOrWhiteSpace(row.Domain) ? null : row.Domain.Trim(),
  165. Potype = string.IsNullOrWhiteSpace(row.Potype) ? null : row.Potype.Trim(),
  166. RawStatus = raw,
  167. CompletionStatus = completion,
  168. SourceSystem = binding.SourceCode,
  169. SourceId = binding.SourceId,
  170. SourceRowId = string.IsNullOrWhiteSpace(row.SourceRowId) ? null : row.SourceRowId.Trim(),
  171. SourceUpdateTime = row.SourceUpdateTime,
  172. SourceUpdateUser = string.IsNullOrWhiteSpace(row.SourceUpdateUser) ? null : row.SourceUpdateUser.Trim(),
  173. });
  174. }
  175. // ④ UPSERT:无条件刷新结论,支持 COMPLETED → NOT_COMPLETED 回退
  176. result.WrittenRows = await UpsertAsync(pending, batchId, now, cancellationToken);
  177. // ⑤ 淘汰:源侧已不存在的行回到 NOT_OBSERVED(本表的表达方式就是「没有这一行」)
  178. // 只在本轮源侧读取成功后执行;任何异常都会在上面抛出并跳过这里。
  179. result.StaleRemoved = await RetireStaleAsync(batchId, tenantId, cancellationToken);
  180. await CompleteRunLogAsync(runLogId, result, now);
  181. return result;
  182. }
  183. catch (Exception ex)
  184. {
  185. _logger.LogError(ex, "[PoCompletion] 批次 {BatchId} 失败", batchId);
  186. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  187. if (!_runLogFinalizer.IsHostStopping)
  188. await FailRunLogAsync(runLogId, ex.Message);
  189. throw;
  190. }
  191. finally
  192. {
  193. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  194. }
  195. }
  196. // ── 源绑定与 fail-closed ────────────────────────────────────────────────────────
  197. /// <summary>
  198. /// 解析并校验源绑定。非开发环境下,未配置、配置为空、或仍指向 DEV/UAT 默认源,
  199. /// 一律直接失败(fail-closed),不得静默继续。
  200. /// </summary>
  201. private async Task<SourceBinding> ResolveSourceBindingAsync(CancellationToken ct)
  202. {
  203. var code = (_opt.SourceCode ?? string.Empty).Trim();
  204. var isDev = _env.IsDevelopment();
  205. if (code.Length == 0)
  206. throw new InvalidOperationException(
  207. "AiDOP:PoCompletion:SourceCode 未配置,无法确定承担 S8_PO_COMPLETION_SOURCE 角色的数据源");
  208. if (!isDev && (_opt.DevOnlySourceCodes ?? Array.Empty<string>())
  209. .Any(x => string.Equals((x ?? string.Empty).Trim(), code, StringComparison.OrdinalIgnoreCase)))
  210. throw new InvalidOperationException(
  211. $"当前环境 {_env.EnvironmentName} 非开发环境,禁止使用 DEV/UAT 数据源 {code};"
  212. + "请在 AiDOP:PoCompletion:SourceCode 显式绑定真实 WMS/MES 源");
  213. var src = await _db.Ado.SqlQuerySingleAsync<SourceBindingRow>(
  214. """
  215. SELECT id AS Id, source_type AS SourceType, IFNULL(db_user,'') AS DbUser
  216. FROM mdp_source
  217. WHERE source_code=@Code AND status=1
  218. LIMIT 1
  219. """,
  220. new List<SugarParameter> { new("@Code", code) });
  221. if (src == null)
  222. throw new InvalidOperationException($"mdp_source 未找到启用源:{code}");
  223. if (!string.Equals(src.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  224. throw new InvalidOperationException($"源 {code} 的 source_type={src.SourceType},不是 DB");
  225. if (string.IsNullOrWhiteSpace(src.DbUser))
  226. throw new InvalidOperationException(
  227. $"源 {code} 未配置数据库账号,连接会回落到主库,不能承担 S8_PO_COMPLETION_SOURCE 角色");
  228. return new SourceBinding { SourceCode = code, SourceId = src.Id };
  229. }
  230. // ── 源侧读取 ────────────────────────────────────────────────────────────────────
  231. /// <summary>
  232. /// 读 Source B 的采购行。范围 = 与自建单推送口径一致的 Potype,
  233. /// 不做无差别全表扫,也不顺带拉无关列。
  234. /// </summary>
  235. private async Task<List<SourceRow>> ReadSourceRowsAsync(ISqlSugarClient remote, CancellationToken ct)
  236. {
  237. var top = _opt.MaxSourceRows > 0 ? _opt.MaxSourceRows : 50000;
  238. var potype = (_opt.SourcePotype ?? string.Empty).Trim();
  239. var sql =
  240. $"""
  241. SELECT TOP {top}
  242. RecID, Domain, Potype, PurOrd, Line, Status, UpdateTime, UpdateUser
  243. FROM PurOrdDetail
  244. WHERE (@Potype = '' OR LOWER(LTRIM(RTRIM(ISNULL(Potype,'')))) = LOWER(@Potype))
  245. """;
  246. var table = await remote.Ado.GetDataTableAsync(sql, new List<SugarParameter> { new("@Potype", potype) });
  247. var list = new List<SourceRow>(table.Rows.Count);
  248. foreach (DataRow r in table.Rows)
  249. {
  250. ct.ThrowIfCancellationRequested();
  251. list.Add(new SourceRow
  252. {
  253. SourceRowId = Str(r, "RecID"),
  254. Domain = Str(r, "Domain"),
  255. Potype = Str(r, "Potype"),
  256. PoNo = Str(r, "PurOrd"),
  257. PoLine = Str(r, "Line"),
  258. RawStatus = Str(r, "Status"),
  259. SourceUpdateTime = Dt(r, "UpdateTime"),
  260. SourceUpdateUser = Str(r, "UpdateUser"),
  261. });
  262. }
  263. return list;
  264. }
  265. // ── 租户归属 ────────────────────────────────────────────────────────────────────
  266. /// <summary>
  267. /// 本地采购单 → 租户候选集。Domain 只用于放宽匹配(自建单本库不填 Domain),
  268. /// <b>绝不</b>用 Domain 反推租户 —— 实测同一个 Domain 8010 横跨多个 Ai-DOP 租户。
  269. /// </summary>
  270. private async Task<Dictionary<string, List<PoTenantCandidate>>> LoadLocalTenantIndexAsync(CancellationToken ct)
  271. {
  272. var rows = await _db.Ado.SqlQueryAsync<PoTenantCandidate>(
  273. """
  274. SELECT PurOrd AS PoNo, IFNULL(Domain,'') AS Domain, tenant_id AS TenantId
  275. FROM PurOrdDetail
  276. WHERE IFNULL(tenant_id,0) > 0
  277. GROUP BY PurOrd, IFNULL(Domain,''), tenant_id
  278. """);
  279. var map = new Dictionary<string, List<PoTenantCandidate>>(StringComparer.OrdinalIgnoreCase);
  280. foreach (var r in rows)
  281. {
  282. if (string.IsNullOrWhiteSpace(r.PoNo)) continue;
  283. var key = r.PoNo.Trim();
  284. if (!map.TryGetValue(key, out var bucket))
  285. map[key] = bucket = new List<PoTenantCandidate>();
  286. bucket.Add(r);
  287. }
  288. return map;
  289. }
  290. /// <summary>
  291. /// 候选租户集合。Domain 只用于<b>放宽</b>匹配(自建单本库不填 Domain,空 Domain 视为通配),
  292. /// <b>绝不</b>反过来用 Domain 推租户 —— 实测同一个 Domain 8010 横跨多个 Ai-DOP 租户。
  293. /// </summary>
  294. public static IEnumerable<long> TenantCandidates(
  295. IEnumerable<PoTenantCandidate> candidates, string poNo, string? domain)
  296. {
  297. var po = (poNo ?? string.Empty).Trim();
  298. var dom = (domain ?? string.Empty).Trim();
  299. return candidates
  300. .Where(x => string.Equals((x.PoNo ?? string.Empty).Trim(), po, StringComparison.OrdinalIgnoreCase))
  301. .Where(x => (x.Domain ?? string.Empty).Trim().Length == 0
  302. || string.Equals((x.Domain ?? string.Empty).Trim(), dom, StringComparison.OrdinalIgnoreCase))
  303. .Select(x => x.TenantId)
  304. .Distinct();
  305. }
  306. /// <summary>
  307. /// 唯一才用:0 个(UNRESOLVED)或 &gt;1 个(AMBIGUOUS)一律返回 0,交由调用方跳过。
  308. /// <b>禁止</b>退化成「取第一个」—— 那会在多租户重名时把事实写到别人名下。
  309. /// </summary>
  310. public static long ResolveTenantId(
  311. IEnumerable<PoTenantCandidate> candidates, string poNo, string? domain)
  312. {
  313. var hits = TenantCandidates(candidates, poNo, domain).Take(2).ToList();
  314. return hits.Count == 1 ? hits[0] : 0;
  315. }
  316. private static long ResolveTenant(
  317. Dictionary<string, List<PoTenantCandidate>> index, string poNo, string? domain)
  318. => index.TryGetValue(poNo.Trim(), out var bucket) ? ResolveTenantId(bucket, poNo, domain) : 0;
  319. private static int CountTenantCandidates(
  320. Dictionary<string, List<PoTenantCandidate>> index, string poNo, string? domain)
  321. => index.TryGetValue(poNo.Trim(), out var bucket)
  322. ? TenantCandidates(bucket, poNo, domain).Count()
  323. : 0;
  324. // ── 本地采购行校验 ──────────────────────────────────────────────────────────────
  325. private async Task<Dictionary<string, int>> LoadLocalPoLineIndexAsync(CancellationToken ct)
  326. {
  327. var rows = await _db.Ado.SqlQueryAsync<PoLineKeyRow>(
  328. "SELECT tenant_id AS TenantId, po_no AS PoNo, po_line AS PoLine FROM mdp_std_purchase_order");
  329. var map = new Dictionary<string, int>(StringComparer.OrdinalIgnoreCase);
  330. foreach (var r in rows)
  331. {
  332. var key = LineKey(r.TenantId, r.PoNo ?? string.Empty, r.PoLine ?? string.Empty);
  333. map[key] = map.TryGetValue(key, out var n) ? n + 1 : 1;
  334. }
  335. return map;
  336. }
  337. private static string LineKey(long tenantId, string poNo, string poLine)
  338. => $"{tenantId}|{poNo.Trim()}|{poLine.Trim()}";
  339. // ── 写入 ────────────────────────────────────────────────────────────────────────
  340. /// <summary>
  341. /// 幂等 UPSERT。结论列一律 <c>VALUES(...)</c> 无条件覆盖——这是支持
  342. /// COMPLETED → NOT_COMPLETED 回退的前提,绝不能改成「仅当新值为 C 才更新」。
  343. /// </summary>
  344. private async Task<int> UpsertAsync(
  345. List<CompletionRow> rows, string batchId, DateTime now, CancellationToken ct)
  346. {
  347. if (rows.Count == 0) return 0;
  348. var allowedKeys = new HashSet<(long TenantId, string Source)>();
  349. foreach (var key in rows.Select(r => (r.TenantId, r.SourceSystem ?? "")).Distinct())
  350. {
  351. if (await _neutralGate.AllowsAsync(key.TenantId, "PURCHASE", key.Item2))
  352. allowedKeys.Add(key);
  353. }
  354. rows = rows.Where(r => allowedKeys.Contains((r.TenantId, r.SourceSystem ?? ""))).ToList();
  355. if (rows.Count == 0) return 0;
  356. var size = _opt.UpsertBatchSize > 0 ? _opt.UpsertBatchSize : 200;
  357. var written = 0;
  358. foreach (var chunk in Chunk(rows, size))
  359. {
  360. ct.ThrowIfCancellationRequested();
  361. var sb = new StringBuilder();
  362. sb.Append("""
  363. INSERT INTO mdp_std_purchase_order_completion
  364. (tenant_id, factory_id, po_no, po_line, domain, potype,
  365. raw_status, completion_status,
  366. source_system, source_id, source_row_id, source_update_time, source_update_user,
  367. observed_at, sync_batch_id, sync_time)
  368. VALUES
  369. """);
  370. var pars = new List<SugarParameter> { new("@Observed", now), new("@Batch", batchId), new("@Now", now) };
  371. for (var i = 0; i < chunk.Count; i++)
  372. {
  373. var r = chunk[i];
  374. if (i > 0) sb.Append(',');
  375. sb.Append($"(@t{i},1,@p{i},@l{i},@d{i},@y{i},@r{i},@c{i},@s{i},@i{i},@x{i},@u{i},@w{i},@Observed,@Batch,@Now)");
  376. pars.Add(new SugarParameter($"@t{i}", r.TenantId));
  377. pars.Add(new SugarParameter($"@p{i}", r.PoNo));
  378. pars.Add(new SugarParameter($"@l{i}", r.PoLine));
  379. pars.Add(new SugarParameter($"@d{i}", (object?)r.Domain ?? DBNull.Value));
  380. pars.Add(new SugarParameter($"@y{i}", (object?)r.Potype ?? DBNull.Value));
  381. pars.Add(new SugarParameter($"@r{i}", (object?)r.RawStatus ?? DBNull.Value));
  382. pars.Add(new SugarParameter($"@c{i}", r.CompletionStatus));
  383. pars.Add(new SugarParameter($"@s{i}", r.SourceSystem));
  384. pars.Add(new SugarParameter($"@i{i}", (object?)r.SourceId ?? DBNull.Value));
  385. pars.Add(new SugarParameter($"@x{i}", (object?)r.SourceRowId ?? DBNull.Value));
  386. pars.Add(new SugarParameter($"@u{i}", (object?)r.SourceUpdateTime ?? DBNull.Value));
  387. pars.Add(new SugarParameter($"@w{i}", (object?)r.SourceUpdateUser ?? DBNull.Value));
  388. }
  389. sb.Append("""
  390. ON DUPLICATE KEY UPDATE
  391. factory_id=VALUES(factory_id), domain=VALUES(domain), potype=VALUES(potype),
  392. raw_status=VALUES(raw_status), completion_status=VALUES(completion_status),
  393. source_system=VALUES(source_system), source_id=VALUES(source_id),
  394. source_row_id=VALUES(source_row_id),
  395. source_update_time=VALUES(source_update_time),
  396. source_update_user=VALUES(source_update_user),
  397. observed_at=VALUES(observed_at), sync_batch_id=VALUES(sync_batch_id),
  398. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  399. """);
  400. await _db.Ado.ExecuteCommandAsync(sb.ToString(), pars);
  401. written += chunk.Count;
  402. }
  403. return written;
  404. }
  405. /// <summary>
  406. /// 淘汰本轮未覆盖的行 —— 它们在源侧已不存在,应回到 NOT_OBSERVED(即本表无此行)。
  407. ///
  408. /// <para>只在本方法被调用时执行,而调用点在源侧读取与写入全部成功之后;
  409. /// 读取失败会在更早处抛出,因此「源侧抓取失败却把旧事实删掉」不可能发生。</para>
  410. ///
  411. /// <para>单租户刷新只淘汰该租户 —— 租户 A 的刷新绝不允许删掉租户 B 的事实。</para>
  412. /// </summary>
  413. private async Task<int> RetireStaleAsync(string batchId, long tenantId, CancellationToken ct)
  414. {
  415. ct.ThrowIfCancellationRequested();
  416. if (tenantId > 0)
  417. {
  418. return await _db.Ado.ExecuteCommandAsync(
  419. """
  420. DELETE FROM mdp_std_purchase_order_completion
  421. WHERE tenant_id=@TenantId AND IFNULL(sync_batch_id,'') <> @BatchId
  422. """,
  423. new List<SugarParameter> { new("@TenantId", tenantId), new("@BatchId", batchId) });
  424. }
  425. return await _db.Ado.ExecuteCommandAsync(
  426. """
  427. DELETE FROM mdp_std_purchase_order_completion
  428. WHERE tenant_id > 0 AND IFNULL(sync_batch_id,'') <> @BatchId
  429. """,
  430. new List<SugarParameter> { new("@BatchId", batchId) });
  431. }
  432. // ── run log ─────────────────────────────────────────────────────────────────────
  433. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType, long tenantId)
  434. {
  435. await _db.Ado.ExecuteCommandAsync(
  436. """
  437. INSERT INTO mdp_transform_run_log
  438. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  439. VALUES (@TenantId, @JobCode, @JobName, @TriggerType, @BatchId, 'RUNNING', @StartTime)
  440. """,
  441. new SugarParameter("@TenantId", tenantId),
  442. new SugarParameter("@JobCode", JobCode),
  443. new SugarParameter("@JobName", JobName),
  444. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  445. new SugarParameter("@BatchId", batchId),
  446. new SugarParameter("@StartTime", startedAt));
  447. return await _db.Ado.GetLongAsync(
  448. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  449. new List<SugarParameter> { new("@BatchId", batchId) });
  450. }
  451. private async Task CompleteRunLogAsync(long runLogId, PurchaseOrderCompletionSyncResult r, DateTime startedAt)
  452. {
  453. var endedAt = DateTime.Now;
  454. // summary_json 是 JSON 列,必须写合法 JSON —— 纯文本会被 MySQL 直接拒绝。
  455. // 这里只放计数与源编码,不放任何连接信息。
  456. var summary = System.Text.Json.JsonSerializer.Serialize(new
  457. {
  458. sourceCode = r.SourceCode,
  459. sourceRows = r.SourceRows,
  460. matchedRows = r.MatchedRows,
  461. writtenRows = r.WrittenRows,
  462. completedRows = r.CompletedRows,
  463. notCompletedRows = r.NotCompletedRows,
  464. unknownStatusRows = r.UnknownStatusRows,
  465. skippedRows = r.SkippedRows,
  466. tenantUnresolvedRows = r.TenantUnresolvedRows,
  467. unmatchedRows = r.UnmatchedRows,
  468. staleRemoved = r.StaleRemoved,
  469. });
  470. await _db.Ado.ExecuteCommandAsync(
  471. """
  472. UPDATE mdp_transform_run_log
  473. SET status='SUCCESS', stage_rows=@SrcRows, standard_rows=@StdRows, end_time=@EndTime,
  474. duration_ms=@Duration, summary_json=@Summary, update_time=CURRENT_TIMESTAMP
  475. WHERE id=@Id
  476. """,
  477. new SugarParameter("@SrcRows", r.SourceRows),
  478. new SugarParameter("@StdRows", r.WrittenRows),
  479. new SugarParameter("@EndTime", endedAt),
  480. new SugarParameter("@Duration", (int)(endedAt - startedAt).TotalMilliseconds),
  481. // 刻意不截断:这是固定形状的计数 JSON(长度有界),截断会产生非法 JSON 并让整轮失败
  482. new SugarParameter("@Summary", summary),
  483. new SugarParameter("@Id", runLogId));
  484. }
  485. private async Task FailRunLogAsync(long runLogId, string message)
  486. {
  487. await _db.Ado.ExecuteCommandAsync(
  488. """
  489. UPDATE mdp_transform_run_log
  490. SET status='FAILED', end_time=@EndTime, error_message=@Msg, update_time=CURRENT_TIMESTAMP
  491. WHERE id=@Id
  492. """,
  493. new SugarParameter("@EndTime", DateTime.Now),
  494. new SugarParameter("@Msg", Truncate(message, 900)),
  495. new SugarParameter("@Id", runLogId));
  496. }
  497. // ── 工具 ────────────────────────────────────────────────────────────────────────
  498. private static string NormalizeTriggerType(string triggerType)
  499. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  500. private static string Truncate(string? s, int max)
  501. => string.IsNullOrEmpty(s) ? string.Empty : (s.Length <= max ? s : s[..max]);
  502. private static string? Str(DataRow row, string col)
  503. {
  504. if (!row.Table.Columns.Contains(col)) return null;
  505. var v = row[col];
  506. return v == null || v == DBNull.Value ? null : Convert.ToString(v);
  507. }
  508. private static DateTime? Dt(DataRow row, string col)
  509. {
  510. if (!row.Table.Columns.Contains(col)) return null;
  511. var v = row[col];
  512. if (v == null || v == DBNull.Value) return null;
  513. return Convert.ToDateTime(v);
  514. }
  515. private static IEnumerable<List<T>> Chunk<T>(List<T> source, int size)
  516. {
  517. for (var i = 0; i < source.Count; i += size)
  518. yield return source.GetRange(i, Math.Min(size, source.Count - i));
  519. }
  520. // ── 内部模型 ────────────────────────────────────────────────────────────────────
  521. private sealed class SourceRow
  522. {
  523. public string? SourceRowId { get; set; }
  524. public string? Domain { get; set; }
  525. public string? Potype { get; set; }
  526. public string PoNo { get; set; } = string.Empty;
  527. public string PoLine { get; set; } = string.Empty;
  528. public string? RawStatus { get; set; }
  529. public DateTime? SourceUpdateTime { get; set; }
  530. public string? SourceUpdateUser { get; set; }
  531. }
  532. private sealed class CompletionRow
  533. {
  534. public long TenantId { get; set; }
  535. public string PoNo { get; set; } = string.Empty;
  536. public string PoLine { get; set; } = string.Empty;
  537. public string? Domain { get; set; }
  538. public string? Potype { get; set; }
  539. public string? RawStatus { get; set; }
  540. public string CompletionStatus { get; set; } = StatusUnknown;
  541. public string SourceSystem { get; set; } = string.Empty;
  542. public long? SourceId { get; set; }
  543. public string? SourceRowId { get; set; }
  544. public DateTime? SourceUpdateTime { get; set; }
  545. public string? SourceUpdateUser { get; set; }
  546. }
  547. private sealed class PoLineKeyRow
  548. {
  549. public long TenantId { get; set; }
  550. public string? PoNo { get; set; }
  551. public string? PoLine { get; set; }
  552. }
  553. private sealed class SourceBinding
  554. {
  555. public string SourceCode { get; set; } = string.Empty;
  556. public long? SourceId { get; set; }
  557. }
  558. private sealed class SourceBindingRow
  559. {
  560. public long Id { get; set; }
  561. public string? SourceType { get; set; }
  562. public string? DbUser { get; set; }
  563. }
  564. }
  565. /// <summary>
  566. /// 本地采购单 → 租户候选。<c>Domain</c> 为空表示自建单(本库不填 165 账套),视为通配。
  567. /// 仅用于把 Source B 的行归到某个 Ai-DOP 租户,<b>不是</b>租户边界本身。
  568. /// </summary>
  569. public sealed class PoTenantCandidate
  570. {
  571. public string PoNo { get; set; } = string.Empty;
  572. public string Domain { get; set; } = string.Empty;
  573. public long TenantId { get; set; }
  574. }
  575. /// <summary>S3 采购完成态投影结果。</summary>
  576. public sealed class PurchaseOrderCompletionSyncResult
  577. {
  578. public string BatchId { get; set; } = string.Empty;
  579. public long RunLogId { get; set; }
  580. /// <summary>承担本次角色的数据源编码(不含任何连接信息)。</summary>
  581. public string SourceCode { get; set; } = string.Empty;
  582. /// <summary>源侧读到的行数。</summary>
  583. public int SourceRows { get; set; }
  584. /// <summary>唯一命中本地采购行的行数。</summary>
  585. public int MatchedRows { get; set; }
  586. /// <summary>实际写入(新增或刷新)的行数。</summary>
  587. public int WrittenRows { get; set; }
  588. /// <summary>被跳过的源行数(租户不可解析 + 本库无对应采购行 + 缺关键字段)。</summary>
  589. public int SkippedRows { get; set; }
  590. /// <summary>其中:租户无法唯一解析。</summary>
  591. public int TenantUnresolvedRows { get; set; }
  592. /// <summary>其中:本库无对应采购行。</summary>
  593. public int UnmatchedRows { get; set; }
  594. public int CompletedRows { get; set; }
  595. public int NotCompletedRows { get; set; }
  596. /// <summary>源侧给了 NULL 或无法识别的状态值。</summary>
  597. public int UnknownStatusRows { get; set; }
  598. /// <summary>本轮源侧已不存在、被回收为 NOT_OBSERVED 的行数。</summary>
  599. public int StaleRemoved { get; set; }
  600. }