MdpDbPushExecutor.cs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306
  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. "ScheduleResultOpMaster", "rf_serialnumber",
  23. // WP7 基础数据 + 任务(DOP 独占)
  24. "LocationMaster", "LocationShelfMaster", "DepartmentMaster",
  25. "EmployeeMaster", "EmpWorkDutyMaster", "LinePrinter",
  26. "MobileTask",
  27. // 单号计数器(NbrSequenceService 直写为主;Outbox 备用)
  28. "NbrDayInfo"
  29. };
  30. private readonly MdpSourceScopeFactory _scopeFactory;
  31. public MdpDbPushExecutor(MdpSourceScopeFactory scopeFactory)
  32. {
  33. _scopeFactory = scopeFactory;
  34. }
  35. public string SupportedType => "DB";
  36. public async Task<MdpPushResult> PushAsync(MdpSource source, MdpOutbox item, CancellationToken ct = default)
  37. {
  38. if (source == null) return MdpPushResult.Fail("source 为空");
  39. if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  40. return MdpPushResult.Fail($"源 {source.SourceCode} 的 source_type={source.SourceType},不是 DB");
  41. DbPushPayload payload;
  42. try
  43. {
  44. payload = ParsePayload(item.PayloadJson);
  45. }
  46. catch (Exception ex)
  47. {
  48. return MdpPushResult.Fail($"payload 解析失败:{ex.Message}");
  49. }
  50. ValidateTable(payload.Table);
  51. ValidateColumns(payload.Keys, "keys");
  52. ValidateColumns(payload.Insert, "insert");
  53. ValidateColumns(payload.Update, "update");
  54. ValidateColumns(payload.Expect, "expect");
  55. var op = (payload.Op ?? "").Trim().ToUpperInvariant();
  56. try
  57. {
  58. var db = await _scopeFactory.GetScopeAsync(source.SourceCode, ct);
  59. return op switch
  60. {
  61. "INSERT" => await ExecInsertAsync(db, payload, ct),
  62. "UPDATE" => await ExecUpdateAsync(db, payload, item.ActionCode, ct),
  63. "UPSERT" => await ExecUpsertAsync(db, payload, item.ActionCode, ct),
  64. _ => MdpPushResult.Fail($"不支持的 op:{payload.Op}")
  65. };
  66. }
  67. catch (Exception ex)
  68. {
  69. return MdpPushResult.Fail($"DB 推送异常:{ex.Message}");
  70. }
  71. }
  72. private static async Task<MdpPushResult> ExecInsertAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  73. {
  74. if (p.Keys.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 keys");
  75. if (p.Insert.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 insert 列");
  76. var arrayErr = RejectArrayValues(p.Keys, "keys")
  77. ?? RejectArrayValues(p.Insert, "insert");
  78. if (arrayErr != null) return MdpPushResult.Fail(arrayErr);
  79. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  80. return MdpPushResult.Skip("{\"idempotent\":true}");
  81. var cols = p.Insert.Keys.ToList();
  82. var colSql = string.Join(",", cols);
  83. var paramSql = string.Join(",", cols.Select((_, i) => $"@v{i}"));
  84. var pars = cols.Select((c, i) => new SugarParameter($"@v{i}", CoerceDbValue(p.Insert[c]) ?? DBNull.Value)).ToArray();
  85. var sql = $"INSERT INTO [{p.Table}] ({colSql}) VALUES ({paramSql})";
  86. var n = await db.Ado.ExecuteCommandAsync(sql, pars);
  87. return MdpPushResult.Ok(n);
  88. }
  89. private static async Task<MdpPushResult> ExecUpdateAsync(
  90. ISqlSugarClient db, DbPushPayload p, string? actionCode, CancellationToken ct)
  91. {
  92. if (p.Keys.Count == 0)
  93. throw new InvalidOperationException("UPDATE 禁止无 keys:keys 为空");
  94. if (p.Update.Count == 0) return MdpPushResult.Fail("UPDATE 必须提供 update 列");
  95. var arrayErr = RejectArrayValues(p.Keys, "keys")
  96. ?? RejectArrayValues(p.Update, "update");
  97. if (arrayErr != null) return MdpPushResult.Fail(arrayErr);
  98. var setParts = p.Update.Keys.Select((c, i) => $"[{c}]=@u{i}").ToList();
  99. var whereParts = p.Keys.Keys.Select((c, i) => $"[{c}]=@k{i}").ToList();
  100. var pars = new List<SugarParameter>();
  101. var ui = 0;
  102. foreach (var kv in p.Update) pars.Add(new SugarParameter($"@u{ui++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  103. var ki = 0;
  104. foreach (var kv in p.Keys) pars.Add(new SugarParameter($"@k{ki++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  105. // P-031 K1-b:expect 标量用 =,数组用 IN;where 与参数共用同一 ei
  106. var ei = 0;
  107. foreach (var kv in p.Expect)
  108. {
  109. if (TryAsExpectList(kv.Value, out var list))
  110. {
  111. if (list!.Count == 0)
  112. return MdpPushResult.Fail("expect 数组为空");
  113. var inNames = new List<string>(list.Count);
  114. for (var j = 0; j < list.Count; j++)
  115. {
  116. var pname = $"@e{ei}_{j}";
  117. inNames.Add(pname);
  118. pars.Add(new SugarParameter(pname, CoerceDbValue(list[j]) ?? DBNull.Value));
  119. }
  120. whereParts.Add($"[{kv.Key}] IN ({string.Join(",", inNames)})");
  121. ei++;
  122. }
  123. else
  124. {
  125. whereParts.Add($"[{kv.Key}]=@e{ei}");
  126. pars.Add(new SugarParameter($"@e{ei}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  127. ei++;
  128. }
  129. }
  130. var sql = $"UPDATE [{p.Table}] SET {string.Join(",", setParts)} WHERE {string.Join(" AND ", whereParts)}";
  131. var n = await db.Ado.ExecuteCommandAsync(sql, pars.ToArray());
  132. if (n == 0)
  133. {
  134. // 作废类:目标行已不存在或 expect 不满足(如已投产)→ 幂等成功,避免死信堆积
  135. if (string.Equals(actionCode, "S2_PSD_DEACTIVATE", StringComparison.OrdinalIgnoreCase))
  136. return MdpPushResult.Skip("{\"idempotent\":true,\"reason\":\"no_row_or_expect\"}");
  137. return MdpPushResult.Fail("UPDATE 影响 0 行(行不存在或 expect 不满足)");
  138. }
  139. return MdpPushResult.Ok(n);
  140. }
  141. private static async Task<MdpPushResult> ExecUpsertAsync(
  142. ISqlSugarClient db, DbPushPayload p, string? actionCode, CancellationToken ct)
  143. {
  144. if (p.Keys.Count == 0) return MdpPushResult.Fail("UPSERT 必须提供 keys");
  145. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  146. return await ExecUpdateAsync(db, p, actionCode, ct);
  147. try
  148. {
  149. return await ExecInsertAsync(db, p, ct);
  150. }
  151. catch (Exception ex)
  152. {
  153. // 并发或唯一索引列宽于 keys 时,INSERT 撞唯一约束 → 收敛为 UPDATE
  154. var err = ex.Message ?? "";
  155. if (err.Contains("唯一", StringComparison.Ordinal)
  156. || err.Contains("UNIQUE", StringComparison.OrdinalIgnoreCase)
  157. || err.Contains("duplicate", StringComparison.OrdinalIgnoreCase)
  158. || err.Contains("2601", StringComparison.Ordinal)
  159. || err.Contains("2627", StringComparison.Ordinal))
  160. {
  161. return await ExecUpdateAsync(db, p, actionCode, ct);
  162. }
  163. throw;
  164. }
  165. }
  166. private static async Task<bool> ExistsAsync(
  167. ISqlSugarClient db, string table, Dictionary<string, object?> keys, CancellationToken ct)
  168. {
  169. var where = string.Join(" AND ", keys.Keys.Select((c, i) => $"[{c}]=@k{i}"));
  170. var pars = keys.Select((kv, i) => new SugarParameter($"@k{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray();
  171. var sql = $"SELECT TOP 1 1 FROM [{table}] WHERE {where}";
  172. var obj = await db.Ado.GetScalarAsync(sql, pars);
  173. return obj != null && obj != DBNull.Value;
  174. }
  175. /// <summary>
  176. /// JSON 反序列化常把日期变成带时区的字符串;SQL Server datetime 转换会失败。
  177. /// 可解析则转为 <see cref="DateTime"/>(取本地/Unspecified)。
  178. /// </summary>
  179. private static object? CoerceDbValue(object? value)
  180. {
  181. if (value == null) return null;
  182. if (value is DateTime or DateTimeOffset or bool or byte or short or int or long or float or double or decimal)
  183. return value is DateTimeOffset dto ? dto.LocalDateTime : value;
  184. if (value is not string s) return value;
  185. if (string.IsNullOrWhiteSpace(s)) return s;
  186. if (DateTimeOffset.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dto2))
  187. return dto2.LocalDateTime;
  188. if (DateTime.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.AssumeLocal, out var dt))
  189. return dt;
  190. return s;
  191. }
  192. /// <summary>P-031 K1-b:非 string 的可枚举视为 expect 多值列表。</summary>
  193. private static bool TryAsExpectList(object? value, out List<object?>? list)
  194. {
  195. list = null;
  196. if (value is null or string) return false;
  197. if (value is not IEnumerable enumerable || value is IDictionary) return false;
  198. list = new List<object?>();
  199. foreach (var item in enumerable)
  200. list.Add(item);
  201. return true;
  202. }
  203. /// <summary>P-031 K1-c:keys/insert/update 禁止数组值。</summary>
  204. private static string? RejectArrayValues(Dictionary<string, object?> cols, string part)
  205. {
  206. foreach (var kv in cols)
  207. {
  208. if (kv.Value is null or string) continue;
  209. if (kv.Value is IEnumerable and not IDictionary)
  210. return $"{part} 不支持数组值,仅 expect 支持";
  211. }
  212. return null;
  213. }
  214. private static void ValidateTable(string table)
  215. {
  216. if (string.IsNullOrWhiteSpace(table) || !ColNameRegex.IsMatch(table))
  217. throw new InvalidOperationException($"非法表名:{table}");
  218. if (!AllowedTables.Contains(table))
  219. throw new InvalidOperationException($"表不在回写白名单:{table}");
  220. }
  221. private static void ValidateColumns(Dictionary<string, object?> cols, string part)
  222. {
  223. foreach (var c in cols.Keys)
  224. {
  225. if (!ColNameRegex.IsMatch(c))
  226. throw new InvalidOperationException($"非法列名({part}):{c}");
  227. }
  228. }
  229. private static DbPushPayload ParsePayload(string? json)
  230. {
  231. using var doc = JsonDocument.Parse(string.IsNullOrWhiteSpace(json) ? "{}" : json);
  232. var root = doc.RootElement;
  233. var p = new DbPushPayload
  234. {
  235. Op = root.TryGetProperty("op", out var op) ? op.GetString() ?? "" : "",
  236. Table = root.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "",
  237. Keys = ReadDict(root, "keys"),
  238. Insert = ReadDict(root, "insert"),
  239. Update = ReadDict(root, "update"),
  240. Expect = ReadDict(root, "expect")
  241. };
  242. return p;
  243. }
  244. private static Dictionary<string, object?> ReadDict(JsonElement root, string name)
  245. {
  246. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  247. if (!root.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Object)
  248. return dict;
  249. foreach (var prop in el.EnumerateObject())
  250. dict[prop.Name] = JsonElementToObject(prop.Value);
  251. return dict;
  252. }
  253. private static object? JsonElementToObject(JsonElement el) => el.ValueKind switch
  254. {
  255. JsonValueKind.Null or JsonValueKind.Undefined => null,
  256. JsonValueKind.String => el.GetString(),
  257. JsonValueKind.True => true,
  258. JsonValueKind.False => false,
  259. JsonValueKind.Number => el.TryGetInt64(out var l) ? l : el.GetDecimal(),
  260. // P-031 K1-a:数组必须保留为列表,否则 expect IN 永远不触发
  261. JsonValueKind.Array => el.EnumerateArray().Select(JsonElementToObject).ToList(),
  262. _ => el.GetRawText()
  263. };
  264. private sealed class DbPushPayload
  265. {
  266. public string Op { get; set; } = "";
  267. public string Table { get; set; } = "";
  268. public Dictionary<string, object?> Keys { get; set; } = new();
  269. public Dictionary<string, object?> Insert { get; set; } = new();
  270. public Dictionary<string, object?> Update { get; set; } = new();
  271. public Dictionary<string, object?> Expect { get; set; } = new();
  272. }
  273. }