MdpDbPushExecutor.cs 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. using System.Globalization;
  2. using System.Text.Json;
  3. using System.Text.RegularExpressions;
  4. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  5. using SqlSugar;
  6. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  7. /// <summary>
  8. /// 字段级 DB 回写:按 payload_json 对 165 SQL Server 执行 INSERT/UPDATE/UPSERT。
  9. /// 只保证「不越表、不无条件更新、不注入」;列级白名单由 WP4 业务编排保证。
  10. /// </summary>
  11. public sealed class MdpDbPushExecutor : IMdpTargetPushExecutor, ITransient
  12. {
  13. private static readonly Regex ColNameRegex = new(@"^[A-Za-z0-9_]+$", RegexOptions.Compiled);
  14. private static readonly HashSet<string> AllowedTables = new(StringComparer.OrdinalIgnoreCase)
  15. {
  16. // WP4 共管表
  17. "WorkOrdMaster", "WorkOrdDetail", "WorkOrdRouting",
  18. "PurOrdMaster", "PurOrdDetail", "PurOrdDetailBatch",
  19. "NbrMaster", "NbrDetail",
  20. "PeriodSequenceDet", "srm_polist_ds",
  21. "ScheduleResultOpMaster", "rf_serialnumber",
  22. // WP7 基础数据 + 任务(DOP 独占)
  23. "LocationMaster", "LocationShelfMaster", "DepartmentMaster",
  24. "EmployeeMaster", "EmpWorkDutyMaster", "LinePrinter",
  25. "MobileTask",
  26. // 单号计数器(NbrSequenceService 直写为主;Outbox 备用)
  27. "NbrDayInfo"
  28. };
  29. private readonly MdpSourceScopeFactory _scopeFactory;
  30. public MdpDbPushExecutor(MdpSourceScopeFactory scopeFactory)
  31. {
  32. _scopeFactory = scopeFactory;
  33. }
  34. public string SupportedType => "DB";
  35. public async Task<MdpPushResult> PushAsync(MdpSource source, MdpOutbox item, CancellationToken ct = default)
  36. {
  37. if (source == null) return MdpPushResult.Fail("source 为空");
  38. if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  39. return MdpPushResult.Fail($"源 {source.SourceCode} 的 source_type={source.SourceType},不是 DB");
  40. DbPushPayload payload;
  41. try
  42. {
  43. payload = ParsePayload(item.PayloadJson);
  44. }
  45. catch (Exception ex)
  46. {
  47. return MdpPushResult.Fail($"payload 解析失败:{ex.Message}");
  48. }
  49. ValidateTable(payload.Table);
  50. ValidateColumns(payload.Keys, "keys");
  51. ValidateColumns(payload.Insert, "insert");
  52. ValidateColumns(payload.Update, "update");
  53. ValidateColumns(payload.Expect, "expect");
  54. var op = (payload.Op ?? "").Trim().ToUpperInvariant();
  55. try
  56. {
  57. var db = await _scopeFactory.GetScopeAsync(source.SourceCode, ct);
  58. return op switch
  59. {
  60. "INSERT" => await ExecInsertAsync(db, payload, ct),
  61. "UPDATE" => await ExecUpdateAsync(db, payload, ct),
  62. "UPSERT" => await ExecUpsertAsync(db, payload, ct),
  63. _ => MdpPushResult.Fail($"不支持的 op:{payload.Op}")
  64. };
  65. }
  66. catch (Exception ex)
  67. {
  68. return MdpPushResult.Fail($"DB 推送异常:{ex.Message}");
  69. }
  70. }
  71. private static async Task<MdpPushResult> ExecInsertAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  72. {
  73. if (p.Keys.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 keys");
  74. if (p.Insert.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 insert 列");
  75. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  76. return MdpPushResult.Skip("{\"idempotent\":true}");
  77. var cols = p.Insert.Keys.ToList();
  78. var colSql = string.Join(",", cols);
  79. var paramSql = string.Join(",", cols.Select((_, i) => $"@v{i}"));
  80. var pars = cols.Select((c, i) => new SugarParameter($"@v{i}", CoerceDbValue(p.Insert[c]) ?? DBNull.Value)).ToArray();
  81. var sql = $"INSERT INTO [{p.Table}] ({colSql}) VALUES ({paramSql})";
  82. var n = await db.Ado.ExecuteCommandAsync(sql, pars);
  83. return MdpPushResult.Ok(n);
  84. }
  85. private static async Task<MdpPushResult> ExecUpdateAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  86. {
  87. if (p.Keys.Count == 0)
  88. throw new InvalidOperationException("UPDATE 禁止无 keys:keys 为空");
  89. if (p.Update.Count == 0) return MdpPushResult.Fail("UPDATE 必须提供 update 列");
  90. var setParts = p.Update.Keys.Select((c, i) => $"[{c}]=@u{i}").ToList();
  91. var whereParts = p.Keys.Keys.Select((c, i) => $"[{c}]=@k{i}").ToList();
  92. whereParts.AddRange(p.Expect.Keys.Select((c, i) => $"[{c}]=@e{i}"));
  93. var pars = new List<SugarParameter>();
  94. var ui = 0;
  95. foreach (var kv in p.Update) pars.Add(new SugarParameter($"@u{ui++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  96. var ki = 0;
  97. foreach (var kv in p.Keys) pars.Add(new SugarParameter($"@k{ki++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  98. var ei = 0;
  99. foreach (var kv in p.Expect) pars.Add(new SugarParameter($"@e{ei++}", CoerceDbValue(kv.Value) ?? DBNull.Value));
  100. var sql = $"UPDATE [{p.Table}] SET {string.Join(",", setParts)} WHERE {string.Join(" AND ", whereParts)}";
  101. var n = await db.Ado.ExecuteCommandAsync(sql, pars.ToArray());
  102. if (n == 0)
  103. return MdpPushResult.Fail("UPDATE 影响 0 行(行不存在或 expect 不满足)");
  104. return MdpPushResult.Ok(n);
  105. }
  106. private static async Task<MdpPushResult> ExecUpsertAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct)
  107. {
  108. if (p.Keys.Count == 0) return MdpPushResult.Fail("UPSERT 必须提供 keys");
  109. if (await ExistsAsync(db, p.Table, p.Keys, ct))
  110. return await ExecUpdateAsync(db, p, ct);
  111. return await ExecInsertAsync(db, p, ct);
  112. }
  113. private static async Task<bool> ExistsAsync(
  114. ISqlSugarClient db, string table, Dictionary<string, object?> keys, CancellationToken ct)
  115. {
  116. var where = string.Join(" AND ", keys.Keys.Select((c, i) => $"[{c}]=@k{i}"));
  117. var pars = keys.Select((kv, i) => new SugarParameter($"@k{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray();
  118. var sql = $"SELECT TOP 1 1 FROM [{table}] WHERE {where}";
  119. var obj = await db.Ado.GetScalarAsync(sql, pars);
  120. return obj != null && obj != DBNull.Value;
  121. }
  122. /// <summary>
  123. /// JSON 反序列化常把日期变成带时区的字符串;SQL Server datetime 转换会失败。
  124. /// 可解析则转为 <see cref="DateTime"/>(取本地/Unspecified)。
  125. /// </summary>
  126. private static object? CoerceDbValue(object? value)
  127. {
  128. if (value == null) return null;
  129. if (value is DateTime or DateTimeOffset or bool or byte or short or int or long or float or double or decimal)
  130. return value is DateTimeOffset dto ? dto.LocalDateTime : value;
  131. if (value is not string s) return value;
  132. if (string.IsNullOrWhiteSpace(s)) return s;
  133. if (DateTimeOffset.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dto2))
  134. return dto2.LocalDateTime;
  135. if (DateTime.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.AssumeLocal, out var dt))
  136. return dt;
  137. return s;
  138. }
  139. private static void ValidateTable(string table)
  140. {
  141. if (string.IsNullOrWhiteSpace(table) || !ColNameRegex.IsMatch(table))
  142. throw new InvalidOperationException($"非法表名:{table}");
  143. if (!AllowedTables.Contains(table))
  144. throw new InvalidOperationException($"表不在回写白名单:{table}");
  145. }
  146. private static void ValidateColumns(Dictionary<string, object?> cols, string part)
  147. {
  148. foreach (var c in cols.Keys)
  149. {
  150. if (!ColNameRegex.IsMatch(c))
  151. throw new InvalidOperationException($"非法列名({part}):{c}");
  152. }
  153. }
  154. private static DbPushPayload ParsePayload(string? json)
  155. {
  156. using var doc = JsonDocument.Parse(string.IsNullOrWhiteSpace(json) ? "{}" : json);
  157. var root = doc.RootElement;
  158. var p = new DbPushPayload
  159. {
  160. Op = root.TryGetProperty("op", out var op) ? op.GetString() ?? "" : "",
  161. Table = root.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "",
  162. Keys = ReadDict(root, "keys"),
  163. Insert = ReadDict(root, "insert"),
  164. Update = ReadDict(root, "update"),
  165. Expect = ReadDict(root, "expect")
  166. };
  167. return p;
  168. }
  169. private static Dictionary<string, object?> ReadDict(JsonElement root, string name)
  170. {
  171. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  172. if (!root.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Object)
  173. return dict;
  174. foreach (var prop in el.EnumerateObject())
  175. dict[prop.Name] = JsonElementToObject(prop.Value);
  176. return dict;
  177. }
  178. private static object? JsonElementToObject(JsonElement el) => el.ValueKind switch
  179. {
  180. JsonValueKind.Null or JsonValueKind.Undefined => null,
  181. JsonValueKind.String => el.GetString(),
  182. JsonValueKind.True => true,
  183. JsonValueKind.False => false,
  184. JsonValueKind.Number => el.TryGetInt64(out var l) ? l : el.GetDecimal(),
  185. _ => el.GetRawText()
  186. };
  187. private sealed class DbPushPayload
  188. {
  189. public string Op { get; set; } = "";
  190. public string Table { get; set; } = "";
  191. public Dictionary<string, object?> Keys { get; set; } = new();
  192. public Dictionary<string, object?> Insert { get; set; } = new();
  193. public Dictionary<string, object?> Update { get; set; } = new();
  194. public Dictionary<string, object?> Expect { get; set; } = new();
  195. }
  196. }