MdpDbPushExecutor.cs 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383
  1. using System.Collections;
  2. using System.Globalization;
  3. using System.Text.Json;
  4. using System.Text.RegularExpressions;
  5. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  6. using SqlSugar;
  7. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  8. /// <summary>
  9. /// 字段级 DB 回写:按 payload_json 对 165 SQL Server 执行 INSERT/UPDATE/UPSERT。
  10. /// 只保证「不越表、不无条件更新、不注入」;列级白名单由 WP4 业务编排保证。
  11. /// </summary>
  12. public sealed class MdpDbPushExecutor : IMdpTargetPushExecutor, ITransient
  13. {
  14. private static readonly Regex ColNameRegex = new(@"^[A-Za-z0-9_]+$", RegexOptions.Compiled);
  15. private static readonly HashSet<string> AllowedTables = new(StringComparer.OrdinalIgnoreCase)
  16. {
  17. // WP4 共管表
  18. "WorkOrdMaster", "WorkOrdDetail", "WorkOrdRouting",
  19. "PurOrdMaster", "PurOrdDetail", "PurOrdDetailBatch",
  20. "NbrMaster", "NbrDetail",
  21. "PeriodSequenceDet", "srm_polist_ds",
  22. "scm_shd", "scm_shdzb", "scm_shdshph", "MissedPrint",
  23. "ScheduleResultOpMaster", "rf_serialnumber",
  24. // WP7 基础数据 + 任务(DOP 独占)
  25. "LocationMaster", "LocationShelfMaster", "DepartmentMaster",
  26. "EmployeeMaster", "EmpWorkDutyMaster", "LinePrinter",
  27. "MobileTask",
  28. // 单号计数器(NbrSequenceService 直写为主;Outbox 备用)
  29. "NbrDayInfo"
  30. };
  31. private readonly MdpSourceScopeFactory _scopeFactory;
  32. public MdpDbPushExecutor(MdpSourceScopeFactory scopeFactory)
  33. {
  34. _scopeFactory = scopeFactory;
  35. }
  36. public string SupportedType => "DB";
  37. public async Task<MdpPushResult> PushAsync(MdpSource source, MdpOutbox item, CancellationToken ct = default)
  38. {
  39. if (source == null) return MdpPushResult.Fail("source 为空");
  40. if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  41. return MdpPushResult.Fail($"源 {source.SourceCode} 的 source_type={source.SourceType},不是 DB");
  42. DbPushPayload payload;
  43. try
  44. {
  45. payload = ParsePayload(item.PayloadJson);
  46. }
  47. catch (Exception ex)
  48. {
  49. return MdpPushResult.Fail($"payload 解析失败:{ex.Message}");
  50. }
  51. ValidateTable(payload.Table);
  52. ValidateColumns(payload.Keys, "keys");
  53. ValidateColumns(payload.Insert, "insert");
  54. ValidateColumns(payload.Update, "update");
  55. ValidateColumns(payload.Expect, "expect");
  56. ValidateResolves(payload.Resolves);
  57. var op = (payload.Op ?? "").Trim().ToUpperInvariant();
  58. try
  59. {
  60. var db = await _scopeFactory.GetScopeAsync(source.SourceCode, ct);
  61. if (payload.Resolves.Count > 0)
  62. {
  63. var resolveErr = await ApplyResolvesAsync(db, payload, ct);
  64. if (resolveErr != null) return MdpPushResult.Fail(resolveErr);
  65. }
  66. return op switch
  67. {
  68. "INSERT" => await ExecInsertAsync(db, payload, ct),
  69. "UPDATE" => await ExecUpdateAsync(db, payload, item.ActionCode, ct),
  70. "UPSERT" => await ExecUpsertAsync(db, payload, item.ActionCode, ct),
  71. _ => MdpPushResult.Fail($"不支持的 op:{payload.Op}")
  72. };
  73. }
  74. catch (Exception ex)
  75. {
  76. return MdpPushResult.Fail($"DB 推送异常:{ex.Message}");
  77. }
  78. }
  79. private static async Task<MdpPushResult> ExecInsertAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  80. {
  81. if (p.Keys.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 keys");
  82. if (p.Insert.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 insert 列");
  83. var arrayErr = RejectArrayValues(p.Keys, "keys")
  84. ?? RejectArrayValues(p.Insert, "insert");
  85. if (arrayErr != null) return MdpPushResult.Fail(arrayErr);
  86. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  87. return MdpPushResult.Skip("{\"idempotent\":true}");
  88. var cols = p.Insert.Keys.ToList();
  89. var colSql = string.Join(",", cols);
  90. var paramSql = string.Join(",", cols.Select((_, i) => $"@v{i}"));
  91. var pars = cols.Select((c, i) => new SugarParameter($"@v{i}", CoerceDbValue(p.Insert[c]) ?? DBNull.Value)).ToArray();
  92. var sql = $"INSERT INTO [{p.Table}] ({colSql}) VALUES ({paramSql})";
  93. var n = await db.Ado.ExecuteCommandAsync(sql, pars);
  94. return MdpPushResult.Ok(n);
  95. }
  96. private static async Task<MdpPushResult> ExecUpdateAsync(
  97. ISqlSugarClient db, DbPushPayload p, string? actionCode, CancellationToken ct)
  98. {
  99. if (p.Keys.Count == 0)
  100. throw new InvalidOperationException("UPDATE 禁止无 keys:keys 为空");
  101. if (p.Update.Count == 0) return MdpPushResult.Fail("UPDATE 必须提供 update 列");
  102. var arrayErr = RejectArrayValues(p.Keys, "keys")
  103. ?? RejectArrayValues(p.Update, "update");
  104. if (arrayErr != null) return MdpPushResult.Fail(arrayErr);
  105. var setParts = p.Update.Keys.Select((c, i) => $"[{c}]=@u{i}").ToList();
  106. var whereParts = p.Keys.Keys.Select((c, i) => $"[{c}]=@k{i}").ToList();
  107. var pars = new List<SugarParameter>();
  108. var ui = 0;
  109. foreach (var kv in p.Update) pars.Add(new SugarParameter($"@u{ui++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  110. var ki = 0;
  111. foreach (var kv in p.Keys) pars.Add(new SugarParameter($"@k{ki++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  112. // P-031 K1-b:expect 标量用 =,数组用 IN;where 与参数共用同一 ei
  113. var ei = 0;
  114. foreach (var kv in p.Expect)
  115. {
  116. if (TryAsExpectList(kv.Value, out var list))
  117. {
  118. if (list!.Count == 0)
  119. return MdpPushResult.Fail("expect 数组为空");
  120. var inNames = new List<string>(list.Count);
  121. for (var j = 0; j < list.Count; j++)
  122. {
  123. var pname = $"@e{ei}_{j}";
  124. inNames.Add(pname);
  125. pars.Add(new SugarParameter(pname, CoerceDbValue(list[j]) ?? DBNull.Value));
  126. }
  127. whereParts.Add($"[{kv.Key}] IN ({string.Join(",", inNames)})");
  128. ei++;
  129. }
  130. else
  131. {
  132. whereParts.Add($"[{kv.Key}]=@e{ei}");
  133. pars.Add(new SugarParameter($"@e{ei}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  134. ei++;
  135. }
  136. }
  137. var sql = $"UPDATE [{p.Table}] SET {string.Join(",", setParts)} WHERE {string.Join(" AND ", whereParts)}";
  138. var n = await db.Ado.ExecuteCommandAsync(sql, pars.ToArray());
  139. if (n == 0)
  140. {
  141. // 作废类:目标行已不存在或 expect 不满足(如已投产)→ 幂等成功,避免死信堆积
  142. if (string.Equals(actionCode, "S2_PSD_DEACTIVATE", StringComparison.OrdinalIgnoreCase))
  143. return MdpPushResult.Skip("{\"idempotent\":true,\"reason\":\"no_row_or_expect\"}");
  144. return MdpPushResult.Fail("UPDATE 影响 0 行(行不存在或 expect 不满足)");
  145. }
  146. return MdpPushResult.Ok(n);
  147. }
  148. private static async Task<MdpPushResult> ExecUpsertAsync(
  149. ISqlSugarClient db, DbPushPayload p, string? actionCode, CancellationToken ct)
  150. {
  151. if (p.Keys.Count == 0) return MdpPushResult.Fail("UPSERT 必须提供 keys");
  152. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  153. return await ExecUpdateAsync(db, p, actionCode, ct);
  154. try
  155. {
  156. return await ExecInsertAsync(db, p, ct);
  157. }
  158. catch (Exception ex)
  159. {
  160. // 并发或唯一索引列宽于 keys 时,INSERT 撞唯一约束 → 收敛为 UPDATE
  161. var err = ex.Message ?? "";
  162. if (err.Contains("唯一", StringComparison.Ordinal)
  163. || err.Contains("UNIQUE", StringComparison.OrdinalIgnoreCase)
  164. || err.Contains("duplicate", StringComparison.OrdinalIgnoreCase)
  165. || err.Contains("2601", StringComparison.Ordinal)
  166. || err.Contains("2627", StringComparison.Ordinal))
  167. {
  168. return await ExecUpdateAsync(db, p, actionCode, ct);
  169. }
  170. throw;
  171. }
  172. }
  173. private static async Task<bool> ExistsAsync(
  174. ISqlSugarClient db, string table, Dictionary<string, object?> keys, CancellationToken ct)
  175. {
  176. var where = string.Join(" AND ", keys.Keys.Select((c, i) => $"[{c}]=@k{i}"));
  177. var pars = keys.Select((kv, i) => new SugarParameter($"@k{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray();
  178. var sql = $"SELECT TOP 1 1 FROM [{table}] WHERE {where}";
  179. var obj = await db.Ado.GetScalarAsync(sql, pars);
  180. return obj != null && obj != DBNull.Value;
  181. }
  182. /// <summary>
  183. /// JSON 反序列化常把日期变成带时区的字符串;SQL Server datetime 转换会失败。
  184. /// 可解析则转为 <see cref="DateTime"/>(取本地/Unspecified)。
  185. /// </summary>
  186. private static object? CoerceDbValue(object? value)
  187. {
  188. if (value == null) return null;
  189. if (value is DateTime or DateTimeOffset or bool or byte or short or int or long or float or double or decimal)
  190. return value is DateTimeOffset dto ? dto.LocalDateTime : value;
  191. if (value is not string s) return value;
  192. if (string.IsNullOrWhiteSpace(s)) return s;
  193. if (DateTimeOffset.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dto2))
  194. return dto2.LocalDateTime;
  195. if (DateTime.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.AssumeLocal, out var dt))
  196. return dt;
  197. return s;
  198. }
  199. /// <summary>P-031 K1-b:非 string 的可枚举视为 expect 多值列表。</summary>
  200. private static bool TryAsExpectList(object? value, out List<object?>? list)
  201. {
  202. list = null;
  203. if (value is null or string) return false;
  204. if (value is not IEnumerable enumerable || value is IDictionary) return false;
  205. list = new List<object?>();
  206. foreach (var item in enumerable)
  207. list.Add(item);
  208. return true;
  209. }
  210. /// <summary>P-031 K1-c:keys/insert/update 禁止数组值。</summary>
  211. private static string? RejectArrayValues(Dictionary<string, object?> cols, string part)
  212. {
  213. foreach (var kv in cols)
  214. {
  215. if (kv.Value is null or string) continue;
  216. if (kv.Value is IEnumerable and not IDictionary)
  217. return $"{part} 不支持数组值,仅 expect 支持";
  218. }
  219. return null;
  220. }
  221. /// <summary>
  222. /// 目标库自增主键不能由本库 RecID 顶替,外键列须在 165 上现查现填。
  223. /// 上游行还没到位时返回失败,交给 dispatcher 退避重试,等主表推成功后自然收敛。
  224. /// </summary>
  225. private static async Task<string?> ApplyResolvesAsync(
  226. ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  227. {
  228. foreach (var r in p.Resolves)
  229. {
  230. var where = string.Join(" AND ", r.Match.Keys.Select((c, i) => $"[{c}]=@r{i}"));
  231. var pars = r.Match.Select((kv, i) => new SugarParameter($"@r{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray();
  232. var sql = $"SELECT TOP 1 [{r.Column}] FROM [{r.Table}] WHERE {where}";
  233. var val = await db.Ado.GetScalarAsync(sql, pars);
  234. if (val == null || val == DBNull.Value)
  235. return $"resolve 未命中:{r.Table}.{r.Column}(目标列 {r.TargetColumn}),上游行未就绪";
  236. p.Insert[r.TargetColumn] = val;
  237. }
  238. return null;
  239. }
  240. private static void ValidateResolves(List<ResolveSpec> resolves)
  241. {
  242. foreach (var r in resolves)
  243. {
  244. if (!ColNameRegex.IsMatch(r.TargetColumn))
  245. throw new InvalidOperationException($"非法列名(resolve):{r.TargetColumn}");
  246. ValidateTable(r.Table);
  247. if (!ColNameRegex.IsMatch(r.Column))
  248. throw new InvalidOperationException($"非法列名(resolve.column):{r.Column}");
  249. ValidateColumns(r.Match, "resolve.match");
  250. if (r.Match.Count == 0)
  251. throw new InvalidOperationException($"resolve.{r.TargetColumn} 缺少 match,禁止无条件取值");
  252. if (RejectArrayValues(r.Match, "resolve.match") is { } err)
  253. throw new InvalidOperationException(err);
  254. }
  255. }
  256. private static void ValidateTable(string table)
  257. {
  258. if (string.IsNullOrWhiteSpace(table) || !ColNameRegex.IsMatch(table))
  259. throw new InvalidOperationException($"非法表名:{table}");
  260. if (!AllowedTables.Contains(table))
  261. throw new InvalidOperationException($"表不在回写白名单:{table}");
  262. }
  263. private static void ValidateColumns(Dictionary<string, object?> cols, string part)
  264. {
  265. foreach (var c in cols.Keys)
  266. {
  267. if (!ColNameRegex.IsMatch(c))
  268. throw new InvalidOperationException($"非法列名({part}):{c}");
  269. }
  270. }
  271. private static DbPushPayload ParsePayload(string? json)
  272. {
  273. using var doc = JsonDocument.Parse(string.IsNullOrWhiteSpace(json) ? "{}" : json);
  274. var root = doc.RootElement;
  275. var p = new DbPushPayload
  276. {
  277. Op = root.TryGetProperty("op", out var op) ? op.GetString() ?? "" : "",
  278. Table = root.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "",
  279. Keys = ReadDict(root, "keys"),
  280. Insert = ReadDict(root, "insert"),
  281. Update = ReadDict(root, "update"),
  282. Expect = ReadDict(root, "expect"),
  283. Resolves = ReadResolves(root)
  284. };
  285. return p;
  286. }
  287. private static List<ResolveSpec> ReadResolves(JsonElement root)
  288. {
  289. var list = new List<ResolveSpec>();
  290. if (!root.TryGetProperty("resolve", out var el) || el.ValueKind != JsonValueKind.Object)
  291. return list;
  292. foreach (var prop in el.EnumerateObject())
  293. {
  294. if (prop.Value.ValueKind != JsonValueKind.Object)
  295. throw new InvalidOperationException($"resolve.{prop.Name} 必须是对象");
  296. var spec = new ResolveSpec
  297. {
  298. TargetColumn = prop.Name,
  299. Table = prop.Value.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "",
  300. Column = prop.Value.TryGetProperty("column", out var c) ? c.GetString() ?? "" : "",
  301. Match = ReadDict(prop.Value, "match")
  302. };
  303. list.Add(spec);
  304. }
  305. return list;
  306. }
  307. private static Dictionary<string, object?> ReadDict(JsonElement root, string name)
  308. {
  309. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  310. if (!root.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Object)
  311. return dict;
  312. foreach (var prop in el.EnumerateObject())
  313. dict[prop.Name] = JsonElementToObject(prop.Value);
  314. return dict;
  315. }
  316. private static object? JsonElementToObject(JsonElement el) => el.ValueKind switch
  317. {
  318. JsonValueKind.Null or JsonValueKind.Undefined => null,
  319. JsonValueKind.String => el.GetString(),
  320. JsonValueKind.True => true,
  321. JsonValueKind.False => false,
  322. JsonValueKind.Number => el.TryGetInt64(out var l) ? l : el.GetDecimal(),
  323. // P-031 K1-a:数组必须保留为列表,否则 expect IN 永远不触发
  324. JsonValueKind.Array => el.EnumerateArray().Select(JsonElementToObject).ToList(),
  325. _ => el.GetRawText()
  326. };
  327. private sealed class DbPushPayload
  328. {
  329. public string Op { get; set; } = "";
  330. public string Table { get; set; } = "";
  331. public Dictionary<string, object?> Keys { get; set; } = new();
  332. public Dictionary<string, object?> Insert { get; set; } = new();
  333. public Dictionary<string, object?> Update { get; set; } = new();
  334. public Dictionary<string, object?> Expect { get; set; } = new();
  335. public List<ResolveSpec> Resolves { get; set; } = new();
  336. }
  337. private sealed class ResolveSpec
  338. {
  339. public string TargetColumn { get; set; } = "";
  340. public string Table { get; set; } = "";
  341. public string Column { get; set; } = "";
  342. public Dictionary<string, object?> Match { get; set; } = new();
  343. }
  344. }