MdpDbPushExecutor.cs 25 KB

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