MdpDbPushExecutor.cs 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562
  1. using System.Collections;
  2. using System.Globalization;
  3. using System.Text;
  4. using System.Text.Json;
  5. using System.Text.RegularExpressions;
  6. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  7. using SqlSugar;
  8. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  9. /// <summary>
  10. /// 字段级 DB 回写:按 payload_json 对 165 SQL Server 执行 INSERT/UPDATE/UPSERT;
  11. /// 另支持受限 <c>op=PROC</c>(过程名白名单,见 <see cref="AllowedProcs"/>)。
  12. /// 只保证「不越表、不无条件更新、不注入」;列级白名单由 WP4 业务编排保证。
  13. /// </summary>
  14. public sealed class MdpDbPushExecutor : IMdpTargetPushExecutor, ITransient
  15. {
  16. private static readonly Regex ColNameRegex = new(@"^[A-Za-z0-9_]+$", RegexOptions.Compiled);
  17. private static readonly HashSet<string> AllowedTables = new(StringComparer.OrdinalIgnoreCase)
  18. {
  19. // WP4 共管表
  20. "WorkOrdMaster", "WorkOrdDetail", "WorkOrdRouting",
  21. "PurOrdMaster", "PurOrdDetail", "PurOrdDetailBatch",
  22. "NbrMaster", "NbrDetail",
  23. "PeriodSequenceDet", "srm_polist_ds",
  24. "scm_shd", "scm_shdzb", "scm_shdshph", "MissedPrint",
  25. "MissedPrintTransHist",
  26. "qms_qcp_insappnentry", "qms_qcp_inspbill",
  27. "ScheduleResultOpMaster", "rf_serialnumber",
  28. // WP7 基础数据 + 任务(DOP 独占)
  29. "LocationMaster", "LocationShelfMaster", "DepartmentMaster",
  30. "EmployeeMaster", "EmpWorkDutyMaster", "LinePrinter",
  31. "MobileTask",
  32. // 单号计数器(NbrSequenceService 直写为主;Outbox 备用)
  33. "NbrDayInfo",
  34. // S5 IQC QcCheck:resolve 读不良品库位(只读,不写)
  35. "PurOrdControl"
  36. };
  37. /// <summary>
  38. /// op=PROC 过程名白名单。本轮仅放 <c>pr_WMS_BPM_SaveInvUpShelf</c>。
  39. /// 依据:doc/plan/S5-IQC合格上架过账(QcCheck)执行任务书.md §3 Q9;白名单变更须同级评审。
  40. /// </summary>
  41. private static readonly HashSet<string> AllowedProcs = new(StringComparer.OrdinalIgnoreCase)
  42. {
  43. "pr_WMS_BPM_SaveInvUpShelf"
  44. };
  45. private readonly MdpSourceScopeFactory _scopeFactory;
  46. public MdpDbPushExecutor(MdpSourceScopeFactory scopeFactory)
  47. {
  48. _scopeFactory = scopeFactory;
  49. }
  50. public string SupportedType => "DB";
  51. public async Task<MdpPushResult> PushAsync(MdpSource source, MdpOutbox item, CancellationToken ct = default)
  52. {
  53. if (source == null) return MdpPushResult.Fail("source 为空");
  54. if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  55. return MdpPushResult.Fail($"源 {source.SourceCode} 的 source_type={source.SourceType},不是 DB");
  56. DbPushPayload payload;
  57. try
  58. {
  59. payload = ParsePayload(item.PayloadJson);
  60. }
  61. catch (Exception ex)
  62. {
  63. return MdpPushResult.Fail($"payload 解析失败:{ex.Message}");
  64. }
  65. var op = (payload.Op ?? "").Trim().ToUpperInvariant();
  66. try
  67. {
  68. if (op == "PROC")
  69. {
  70. ValidateProc(payload.Proc);
  71. ValidateColumns(payload.Args, "args");
  72. ValidateResolves(payload.Resolves);
  73. if (RejectArrayValues(payload.Args, "args") is { } argsErr)
  74. return MdpPushResult.Fail(argsErr);
  75. }
  76. else
  77. {
  78. ValidateTable(payload.Table);
  79. ValidateColumns(payload.Keys, "keys");
  80. ValidateColumns(payload.Insert, "insert");
  81. ValidateColumns(payload.Update, "update");
  82. ValidateColumns(payload.Expect, "expect");
  83. ValidateResolves(payload.Resolves);
  84. }
  85. var db = await _scopeFactory.GetScopeAsync(source.SourceCode, ct);
  86. if (payload.Resolves.Count > 0)
  87. {
  88. var resolveErr = await ApplyResolvesAsync(db, payload, op, ct);
  89. if (resolveErr != null) return MdpPushResult.Fail(resolveErr);
  90. }
  91. return op switch
  92. {
  93. "INSERT" => await ExecInsertAsync(db, payload, ct),
  94. "UPDATE" => await ExecUpdateAsync(db, payload, item.ActionCode, ct),
  95. "UPSERT" => await ExecUpsertAsync(db, payload, item.ActionCode, ct),
  96. "PROC" => await ExecProcAsync(db, payload, ct),
  97. _ => MdpPushResult.Fail($"不支持的 op:{payload.Op}")
  98. };
  99. }
  100. catch (Exception ex)
  101. {
  102. return MdpPushResult.Fail($"DB 推送异常:{ex.Message}");
  103. }
  104. }
  105. private static async Task<MdpPushResult> ExecInsertAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  106. {
  107. if (p.Keys.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 keys");
  108. if (p.Insert.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 insert 列");
  109. var arrayErr = RejectArrayValues(p.Keys, "keys")
  110. ?? RejectArrayValues(p.Insert, "insert");
  111. if (arrayErr != null) return MdpPushResult.Fail(arrayErr);
  112. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  113. return MdpPushResult.Skip("{\"idempotent\":true}");
  114. var cols = p.Insert.Keys.ToList();
  115. var colSql = string.Join(",", cols);
  116. var paramSql = string.Join(",", cols.Select((_, i) => $"@v{i}"));
  117. var pars = cols.Select((c, i) => new SugarParameter($"@v{i}", CoerceDbValue(p.Insert[c]) ?? DBNull.Value)).ToArray();
  118. var sql = $"INSERT INTO [{p.Table}] ({colSql}) VALUES ({paramSql})";
  119. var n = await db.Ado.ExecuteCommandAsync(sql, pars);
  120. return MdpPushResult.Ok(n);
  121. }
  122. private static async Task<MdpPushResult> ExecUpdateAsync(
  123. ISqlSugarClient db, DbPushPayload p, string? actionCode, CancellationToken ct)
  124. {
  125. if (p.Keys.Count == 0)
  126. throw new InvalidOperationException("UPDATE 禁止无 keys:keys 为空");
  127. if (p.Update.Count == 0) return MdpPushResult.Fail("UPDATE 必须提供 update 列");
  128. var arrayErr = RejectArrayValues(p.Keys, "keys")
  129. ?? RejectArrayValues(p.Update, "update");
  130. if (arrayErr != null) return MdpPushResult.Fail(arrayErr);
  131. var setParts = p.Update.Keys.Select((c, i) => $"[{c}]=@u{i}").ToList();
  132. var whereParts = p.Keys.Keys.Select((c, i) => $"[{c}]=@k{i}").ToList();
  133. var pars = new List<SugarParameter>();
  134. var ui = 0;
  135. foreach (var kv in p.Update) pars.Add(new SugarParameter($"@u{ui++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  136. var ki = 0;
  137. foreach (var kv in p.Keys) pars.Add(new SugarParameter($"@k{ki++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  138. // P-031 K1-b:expect 标量用 =,数组用 IN;where 与参数共用同一 ei
  139. var ei = 0;
  140. foreach (var kv in p.Expect)
  141. {
  142. if (TryAsExpectList(kv.Value, out var list))
  143. {
  144. if (list!.Count == 0)
  145. return MdpPushResult.Fail("expect 数组为空");
  146. var inNames = new List<string>(list.Count);
  147. for (var j = 0; j < list.Count; j++)
  148. {
  149. var pname = $"@e{ei}_{j}";
  150. inNames.Add(pname);
  151. pars.Add(new SugarParameter(pname, CoerceDbValue(list[j]) ?? DBNull.Value));
  152. }
  153. whereParts.Add($"[{kv.Key}] IN ({string.Join(",", inNames)})");
  154. ei++;
  155. }
  156. else
  157. {
  158. whereParts.Add($"[{kv.Key}]=@e{ei}");
  159. pars.Add(new SugarParameter($"@e{ei}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  160. ei++;
  161. }
  162. }
  163. var sql = $"UPDATE [{p.Table}] SET {string.Join(",", setParts)} WHERE {string.Join(" AND ", whereParts)}";
  164. var n = await db.Ado.ExecuteCommandAsync(sql, pars.ToArray());
  165. if (n == 0)
  166. {
  167. // 作废类:目标行已不存在或 expect 不满足(如已投产)→ 幂等成功,避免死信堆积
  168. if (string.Equals(actionCode, "S2_PSD_DEACTIVATE", StringComparison.OrdinalIgnoreCase))
  169. return MdpPushResult.Skip("{\"idempotent\":true,\"reason\":\"no_row_or_expect\"}");
  170. return MdpPushResult.Fail("UPDATE 影响 0 行(行不存在或 expect 不满足)");
  171. }
  172. return MdpPushResult.Ok(n);
  173. }
  174. private static async Task<MdpPushResult> ExecUpsertAsync(
  175. ISqlSugarClient db, DbPushPayload p, string? actionCode, CancellationToken ct)
  176. {
  177. if (p.Keys.Count == 0) return MdpPushResult.Fail("UPSERT 必须提供 keys");
  178. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  179. return await ExecUpdateAsync(db, p, actionCode, ct);
  180. try
  181. {
  182. return await ExecInsertAsync(db, p, ct);
  183. }
  184. catch (Exception ex)
  185. {
  186. // 并发或唯一索引列宽于 keys 时,INSERT 撞唯一约束 → 收敛为 UPDATE
  187. var err = ex.Message ?? "";
  188. if (err.Contains("唯一", StringComparison.Ordinal)
  189. || err.Contains("UNIQUE", StringComparison.OrdinalIgnoreCase)
  190. || err.Contains("duplicate", StringComparison.OrdinalIgnoreCase)
  191. || err.Contains("2601", StringComparison.Ordinal)
  192. || err.Contains("2627", StringComparison.Ordinal))
  193. {
  194. return await ExecUpdateAsync(db, p, actionCode, ct);
  195. }
  196. throw;
  197. }
  198. }
  199. /// <summary>
  200. /// 调 165 白名单存储过程。关键:部分过程在 <c>@IsProcCall=1</c> 时只开事务不提交,
  201. /// 必须由调用方包事务并按 <c>@@TRANCOUNT</c> 收口(见 QcCheck 任务书 §4.2)。
  202. /// </summary>
  203. private static async Task<MdpPushResult> ExecProcAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  204. {
  205. if (string.IsNullOrWhiteSpace(p.Proc))
  206. return MdpPushResult.Fail("PROC 必须提供 proc");
  207. if (p.Args.Count == 0)
  208. return MdpPushResult.Fail("PROC 必须提供 args");
  209. var outputName = string.IsNullOrWhiteSpace(p.Output) ? "ReturnMsg" : p.Output!;
  210. if (!ColNameRegex.IsMatch(outputName))
  211. return MdpPushResult.Fail($"非法 OUTPUT 参数名:{outputName}");
  212. // QcCheck:SP 在 @Shelf 仅库位号(@InvShelf='')时,「已是此货架」校验不命中,会重复过账。
  213. // 入队前按标签当前库位预检,已在目标库位则幂等跳过。
  214. var preSkip = await TrySkipQcCheckIfAlreadyThereAsync(db, p, ct);
  215. if (preSkip != null) return preSkip;
  216. var successNeedle = string.IsNullOrWhiteSpace(p.Success) ? "成功" : p.Success!;
  217. var argKeys = p.Args.Keys.ToList();
  218. var assignParts = new List<string>(argKeys.Count + 1);
  219. var pars = new List<SugarParameter>(argKeys.Count);
  220. for (var i = 0; i < argKeys.Count; i++)
  221. {
  222. var key = argKeys[i];
  223. assignParts.Add($"@{key}=@p{i}");
  224. pars.Add(new SugarParameter($"@p{i}", CoerceDbValue(p.Args[key]) ?? DBNull.Value));
  225. }
  226. assignParts.Add($"@{outputName}=@ret OUTPUT");
  227. var sql = new StringBuilder();
  228. sql.AppendLine("SET XACT_ABORT ON;");
  229. sql.AppendLine("DECLARE @ret varchar(max) = '';");
  230. sql.AppendLine("BEGIN TRAN;");
  231. sql.Append("EXEC [").Append(p.Proc).Append("] ");
  232. sql.Append(string.Join(", ", assignParts)).AppendLine(";");
  233. sql.AppendLine("IF CHARINDEX(N'成功', ISNULL(@ret,'')) = 0");
  234. sql.AppendLine(" WHILE @@TRANCOUNT > 0 ROLLBACK;");
  235. sql.AppendLine("ELSE");
  236. sql.AppendLine(" WHILE @@TRANCOUNT > 0 COMMIT;");
  237. sql.AppendLine("SELECT @ret AS ReturnMsg;");
  238. // 事务收口固定看「成功」(与 SP 文案「保存成功!」对齐);业务成功/幂等再按 payload.success 判定。
  239. var dt = await db.Ado.GetDataTableAsync(sql.ToString(), pars.ToArray());
  240. var retMsg = "";
  241. if (dt is { Rows.Count: > 0 } && dt.Rows[0][0] != DBNull.Value && dt.Rows[0][0] != null)
  242. retMsg = Convert.ToString(dt.Rows[0][0]) ?? "";
  243. var response = JsonSerializer.Serialize(new { returnMsg = retMsg });
  244. if (p.Idempotent.Count > 0
  245. && p.Idempotent.Any(s => !string.IsNullOrEmpty(s)
  246. && retMsg.Contains(s, StringComparison.Ordinal)))
  247. {
  248. return MdpPushResult.Skip(response);
  249. }
  250. if (retMsg.Contains(successNeedle, StringComparison.Ordinal))
  251. return MdpPushResult.Ok(1, response);
  252. return MdpPushResult.Fail(
  253. string.IsNullOrWhiteSpace(retMsg) ? "PROC 未返回成功消息" : retMsg,
  254. response);
  255. }
  256. /// <summary>
  257. /// QcCheck 幂等预检:MissedPrint.Location 已等于目标库位时跳过,避免同库位重复 iss-tr-ins。
  258. /// </summary>
  259. private static async Task<MdpPushResult?> TrySkipQcCheckIfAlreadyThereAsync(
  260. ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  261. {
  262. if (!string.Equals(p.Proc, "pr_WMS_BPM_SaveInvUpShelf", StringComparison.OrdinalIgnoreCase))
  263. return null;
  264. if (!p.Args.TryGetValue("TransCode", out var tc) || tc == null)
  265. return null;
  266. if (!string.Equals(Convert.ToString(tc), "QcCheck", StringComparison.OrdinalIgnoreCase))
  267. return null;
  268. if (!p.Args.TryGetValue("Details", out var detailsObj) || detailsObj == null)
  269. return null;
  270. if (!p.Args.TryGetValue("Shelf", out var shelfObj) || shelfObj == null)
  271. return null;
  272. var details = Convert.ToString(detailsObj, CultureInfo.InvariantCulture)?.Trim() ?? "";
  273. if (string.IsNullOrEmpty(details) || details.Contains('|', StringComparison.Ordinal)
  274. || details.Contains(',', StringComparison.Ordinal))
  275. return null; // 批量细节不做预检,交给 SP
  276. var shelf = Convert.ToString(shelfObj, CultureInfo.InvariantCulture)?.Trim() ?? "";
  277. if (string.IsNullOrEmpty(shelf)) return null;
  278. var splitIdx = shelf.IndexOf(':');
  279. var locationTo = splitIdx > 0 ? shelf[..splitIdx] : shelf;
  280. var curLocObj = await db.Ado.GetScalarAsync(
  281. "SELECT TOP 1 [Location] FROM [MissedPrint] WHERE [RecID]=@id",
  282. new SugarParameter("@id", details));
  283. var curLoc = curLocObj == null || curLocObj == DBNull.Value
  284. ? ""
  285. : Convert.ToString(curLocObj, CultureInfo.InvariantCulture)?.Trim() ?? "";
  286. if (string.IsNullOrEmpty(curLoc)) return null;
  287. if (string.Equals(curLoc, locationTo, StringComparison.OrdinalIgnoreCase))
  288. {
  289. return MdpPushResult.Skip(JsonSerializer.Serialize(new
  290. {
  291. returnMsg = "已是此货架(执行器预检)",
  292. location = curLoc
  293. }));
  294. }
  295. return null;
  296. }
  297. private static async Task<bool> ExistsAsync(
  298. ISqlSugarClient db, string table, Dictionary<string, object?> keys, CancellationToken ct)
  299. {
  300. var where = string.Join(" AND ", keys.Keys.Select((c, i) => $"[{c}]=@k{i}"));
  301. var pars = keys.Select((kv, i) => new SugarParameter($"@k{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray();
  302. var sql = $"SELECT TOP 1 1 FROM [{table}] WHERE {where}";
  303. var obj = await db.Ado.GetScalarAsync(sql, pars);
  304. return obj != null && obj != DBNull.Value;
  305. }
  306. /// <summary>
  307. /// JSON 反序列化常把日期变成带时区的字符串;SQL Server datetime 转换会失败。
  308. /// 可解析则转为 <see cref="DateTime"/>(取本地/Unspecified)。
  309. /// </summary>
  310. private static object? CoerceDbValue(object? value)
  311. {
  312. if (value == null) return null;
  313. if (value is DateTime or DateTimeOffset or bool or byte or short or int or long or float or double or decimal)
  314. return value is DateTimeOffset dto ? dto.LocalDateTime : value;
  315. if (value is not string s) return value;
  316. if (string.IsNullOrWhiteSpace(s)) return s;
  317. if (DateTimeOffset.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dto2))
  318. return dto2.LocalDateTime;
  319. if (DateTime.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.AssumeLocal, out var dt))
  320. return dt;
  321. return s;
  322. }
  323. /// <summary>P-031 K1-b:非 string 的可枚举视为 expect 多值列表。</summary>
  324. private static bool TryAsExpectList(object? value, out List<object?>? list)
  325. {
  326. list = null;
  327. if (value is null or string) return false;
  328. if (value is not IEnumerable enumerable || value is IDictionary) return false;
  329. list = new List<object?>();
  330. foreach (var item in enumerable)
  331. list.Add(item);
  332. return true;
  333. }
  334. /// <summary>P-031 K1-c:keys/insert/update 禁止数组值。</summary>
  335. private static string? RejectArrayValues(Dictionary<string, object?> cols, string part)
  336. {
  337. foreach (var kv in cols)
  338. {
  339. if (kv.Value is null or string) continue;
  340. if (kv.Value is IEnumerable and not IDictionary)
  341. return $"{part} 不支持数组值,仅 expect 支持";
  342. }
  343. return null;
  344. }
  345. /// <summary>
  346. /// 目标库自增主键不能由本库 RecID 顶替,外键列须在 165 上现查现填。
  347. /// 上游行还没到位时返回失败,交给 dispatcher 退避重试,等主表推成功后自然收敛。
  348. /// </summary>
  349. private static async Task<string?> ApplyResolvesAsync(
  350. ISqlSugarClient db, DbPushPayload p, string op, CancellationToken ct)
  351. {
  352. foreach (var r in p.Resolves)
  353. {
  354. var where = string.Join(" AND ", r.Match.Keys.Select((c, i) => $"[{c}]=@r{i}"));
  355. var pars = r.Match.Select((kv, i) => new SugarParameter($"@r{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray();
  356. var sql = $"SELECT TOP 1 [{r.Column}] FROM [{r.Table}] WHERE {where}";
  357. var val = await db.Ado.GetScalarAsync(sql, pars);
  358. if (val == null || val == DBNull.Value)
  359. return $"resolve 未命中:{r.Table}.{r.Column}(目标列 {r.TargetColumn}),上游行未就绪";
  360. // PROC:resolve 写入 args;其余写入 insert(既有语义)
  361. if (op == "PROC")
  362. p.Args[r.TargetColumn] = Convert.ToString(val, CultureInfo.InvariantCulture);
  363. else
  364. p.Insert[r.TargetColumn] = val;
  365. }
  366. return null;
  367. }
  368. private static void ValidateResolves(List<ResolveSpec> resolves)
  369. {
  370. foreach (var r in resolves)
  371. {
  372. if (!ColNameRegex.IsMatch(r.TargetColumn))
  373. throw new InvalidOperationException($"非法列名(resolve):{r.TargetColumn}");
  374. ValidateTable(r.Table);
  375. if (!ColNameRegex.IsMatch(r.Column))
  376. throw new InvalidOperationException($"非法列名(resolve.column):{r.Column}");
  377. ValidateColumns(r.Match, "resolve.match");
  378. if (r.Match.Count == 0)
  379. throw new InvalidOperationException($"resolve.{r.TargetColumn} 缺少 match,禁止无条件取值");
  380. if (RejectArrayValues(r.Match, "resolve.match") is { } err)
  381. throw new InvalidOperationException(err);
  382. }
  383. }
  384. private static void ValidateTable(string table)
  385. {
  386. if (string.IsNullOrWhiteSpace(table) || !ColNameRegex.IsMatch(table))
  387. throw new InvalidOperationException($"非法表名:{table}");
  388. if (!AllowedTables.Contains(table))
  389. throw new InvalidOperationException($"表不在回写白名单:{table}");
  390. }
  391. private static void ValidateProc(string proc)
  392. {
  393. if (string.IsNullOrWhiteSpace(proc) || !ColNameRegex.IsMatch(proc))
  394. throw new InvalidOperationException($"非法过程名:{proc}");
  395. if (!AllowedProcs.Contains(proc))
  396. throw new InvalidOperationException($"过程不在回写白名单:{proc}");
  397. }
  398. private static void ValidateColumns(Dictionary<string, object?> cols, string part)
  399. {
  400. foreach (var c in cols.Keys)
  401. {
  402. if (!ColNameRegex.IsMatch(c))
  403. throw new InvalidOperationException($"非法列名({part}):{c}");
  404. }
  405. }
  406. private static DbPushPayload ParsePayload(string? json)
  407. {
  408. using var doc = JsonDocument.Parse(string.IsNullOrWhiteSpace(json) ? "{}" : json);
  409. var root = doc.RootElement;
  410. var p = new DbPushPayload
  411. {
  412. Op = root.TryGetProperty("op", out var op) ? op.GetString() ?? "" : "",
  413. Table = root.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "",
  414. Proc = root.TryGetProperty("proc", out var proc) ? proc.GetString() ?? "" : "",
  415. Keys = ReadDict(root, "keys"),
  416. Insert = ReadDict(root, "insert"),
  417. Update = ReadDict(root, "update"),
  418. Expect = ReadDict(root, "expect"),
  419. Args = ReadDict(root, "args"),
  420. Resolves = ReadResolves(root),
  421. Output = root.TryGetProperty("output", out var output) ? output.GetString() ?? "" : "",
  422. Success = root.TryGetProperty("success", out var success) ? success.GetString() ?? "" : "",
  423. Idempotent = ReadStringList(root, "idempotent")
  424. };
  425. return p;
  426. }
  427. private static List<string> ReadStringList(JsonElement root, string name)
  428. {
  429. var list = new List<string>();
  430. if (!root.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Array)
  431. return list;
  432. foreach (var item in el.EnumerateArray())
  433. {
  434. if (item.ValueKind == JsonValueKind.String)
  435. {
  436. var s = item.GetString();
  437. if (!string.IsNullOrEmpty(s)) list.Add(s);
  438. }
  439. }
  440. return list;
  441. }
  442. private static List<ResolveSpec> ReadResolves(JsonElement root)
  443. {
  444. var list = new List<ResolveSpec>();
  445. if (!root.TryGetProperty("resolve", out var el) || el.ValueKind != JsonValueKind.Object)
  446. return list;
  447. foreach (var prop in el.EnumerateObject())
  448. {
  449. if (prop.Value.ValueKind != JsonValueKind.Object)
  450. throw new InvalidOperationException($"resolve.{prop.Name} 必须是对象");
  451. var spec = new ResolveSpec
  452. {
  453. TargetColumn = prop.Name,
  454. Table = prop.Value.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "",
  455. Column = prop.Value.TryGetProperty("column", out var c) ? c.GetString() ?? "" : "",
  456. Match = ReadDict(prop.Value, "match")
  457. };
  458. list.Add(spec);
  459. }
  460. return list;
  461. }
  462. private static Dictionary<string, object?> ReadDict(JsonElement root, string name)
  463. {
  464. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  465. if (!root.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Object)
  466. return dict;
  467. foreach (var prop in el.EnumerateObject())
  468. dict[prop.Name] = JsonElementToObject(prop.Value);
  469. return dict;
  470. }
  471. private static object? JsonElementToObject(JsonElement el) => el.ValueKind switch
  472. {
  473. JsonValueKind.Null or JsonValueKind.Undefined => null,
  474. JsonValueKind.String => el.GetString(),
  475. JsonValueKind.True => true,
  476. JsonValueKind.False => false,
  477. JsonValueKind.Number => el.TryGetInt64(out var l) ? l : el.GetDecimal(),
  478. // P-031 K1-a:数组必须保留为列表,否则 expect IN 永远不触发
  479. JsonValueKind.Array => el.EnumerateArray().Select(JsonElementToObject).ToList(),
  480. _ => el.GetRawText()
  481. };
  482. private sealed class DbPushPayload
  483. {
  484. public string Op { get; set; } = "";
  485. public string Table { get; set; } = "";
  486. public string Proc { get; set; } = "";
  487. public Dictionary<string, object?> Keys { get; set; } = new();
  488. public Dictionary<string, object?> Insert { get; set; } = new();
  489. public Dictionary<string, object?> Update { get; set; } = new();
  490. public Dictionary<string, object?> Expect { get; set; } = new();
  491. public Dictionary<string, object?> Args { get; set; } = new();
  492. public List<ResolveSpec> Resolves { get; set; } = new();
  493. public string Output { get; set; } = "";
  494. public string Success { get; set; } = "";
  495. public List<string> Idempotent { get; set; } = new();
  496. }
  497. private sealed class ResolveSpec
  498. {
  499. public string TargetColumn { get; set; } = "";
  500. public string Table { get; set; } = "";
  501. public string Column { get; set; } = "";
  502. public Dictionary<string, object?> Match { get; set; } = new();
  503. }
  504. }