using System.Globalization; using System.Text.Json; using System.Text.RegularExpressions; using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors; /// /// 字段级 DB 回写:按 payload_json 对 165 SQL Server 执行 INSERT/UPDATE/UPSERT。 /// 只保证「不越表、不无条件更新、不注入」;列级白名单由 WP4 业务编排保证。 /// public sealed class MdpDbPushExecutor : IMdpTargetPushExecutor, ITransient { private static readonly Regex ColNameRegex = new(@"^[A-Za-z0-9_]+$", RegexOptions.Compiled); private static readonly HashSet AllowedTables = new(StringComparer.OrdinalIgnoreCase) { // WP4 共管表 "WorkOrdMaster", "WorkOrdDetail", "WorkOrdRouting", "PurOrdMaster", "PurOrdDetail", "PurOrdDetailBatch", "NbrMaster", "NbrDetail", "PeriodSequenceDet", "srm_polist_ds", "ScheduleResultOpMaster", "rf_serialnumber", // WP7 基础数据 + 任务(DOP 独占) "LocationMaster", "LocationShelfMaster", "DepartmentMaster", "EmployeeMaster", "EmpWorkDutyMaster", "LinePrinter", "MobileTask", // 单号计数器(NbrSequenceService 直写为主;Outbox 备用) "NbrDayInfo" }; private readonly MdpSourceScopeFactory _scopeFactory; public MdpDbPushExecutor(MdpSourceScopeFactory scopeFactory) { _scopeFactory = scopeFactory; } public string SupportedType => "DB"; public async Task PushAsync(MdpSource source, MdpOutbox item, CancellationToken ct = default) { if (source == null) return MdpPushResult.Fail("source 为空"); if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase)) return MdpPushResult.Fail($"源 {source.SourceCode} 的 source_type={source.SourceType},不是 DB"); DbPushPayload payload; try { payload = ParsePayload(item.PayloadJson); } catch (Exception ex) { return MdpPushResult.Fail($"payload 解析失败:{ex.Message}"); } ValidateTable(payload.Table); ValidateColumns(payload.Keys, "keys"); ValidateColumns(payload.Insert, "insert"); ValidateColumns(payload.Update, "update"); ValidateColumns(payload.Expect, "expect"); var op = (payload.Op ?? "").Trim().ToUpperInvariant(); try { var db = await _scopeFactory.GetScopeAsync(source.SourceCode, ct); return op switch { "INSERT" => await ExecInsertAsync(db, payload, ct), "UPDATE" => await ExecUpdateAsync(db, payload, ct), "UPSERT" => await ExecUpsertAsync(db, payload, ct), _ => MdpPushResult.Fail($"不支持的 op:{payload.Op}") }; } catch (Exception ex) { return MdpPushResult.Fail($"DB 推送异常:{ex.Message}"); } } private static async Task ExecInsertAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct) { if (p.Keys.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 keys"); if (p.Insert.Count == 0) return MdpPushResult.Fail("INSERT 必须提供 insert 列"); if (await ExistsAsync(db, p.Table, p.Keys, ct)) return MdpPushResult.Skip("{\"idempotent\":true}"); var cols = p.Insert.Keys.ToList(); var colSql = string.Join(",", cols); var paramSql = string.Join(",", cols.Select((_, i) => $"@v{i}")); var pars = cols.Select((c, i) => new SugarParameter($"@v{i}", CoerceDbValue(p.Insert[c]) ?? DBNull.Value)).ToArray(); var sql = $"INSERT INTO [{p.Table}] ({colSql}) VALUES ({paramSql})"; var n = await db.Ado.ExecuteCommandAsync(sql, pars); return MdpPushResult.Ok(n); } private static async Task ExecUpdateAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct) { if (p.Keys.Count == 0) throw new InvalidOperationException("UPDATE 禁止无 keys:keys 为空"); if (p.Update.Count == 0) return MdpPushResult.Fail("UPDATE 必须提供 update 列"); var setParts = p.Update.Keys.Select((c, i) => $"[{c}]=@u{i}").ToList(); var whereParts = p.Keys.Keys.Select((c, i) => $"[{c}]=@k{i}").ToList(); whereParts.AddRange(p.Expect.Keys.Select((c, i) => $"[{c}]=@e{i}")); var pars = new List(); var ui = 0; foreach (var kv in p.Update) pars.Add(new SugarParameter($"@u{ui++}", CoerceDbValue(kv.Value) ?? DBNull.Value)); var ki = 0; foreach (var kv in p.Keys) pars.Add(new SugarParameter($"@k{ki++}", CoerceDbValue(kv.Value) ?? DBNull.Value)); var ei = 0; foreach (var kv in p.Expect) pars.Add(new SugarParameter($"@e{ei++}", CoerceDbValue(kv.Value) ?? DBNull.Value)); var sql = $"UPDATE [{p.Table}] SET {string.Join(",", setParts)} WHERE {string.Join(" AND ", whereParts)}"; var n = await db.Ado.ExecuteCommandAsync(sql, pars.ToArray()); if (n == 0) return MdpPushResult.Fail("UPDATE 影响 0 行(行不存在或 expect 不满足)"); return MdpPushResult.Ok(n); } private static async Task ExecUpsertAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct) { if (p.Keys.Count == 0) return MdpPushResult.Fail("UPSERT 必须提供 keys"); if (await ExistsAsync(db, p.Table, p.Keys, ct)) return await ExecUpdateAsync(db, p, ct); try { return await ExecInsertAsync(db, p, ct); } catch (Exception ex) { // 并发或唯一索引列宽于 keys 时,INSERT 撞唯一约束 → 收敛为 UPDATE var err = ex.Message ?? ""; if (err.Contains("唯一", StringComparison.Ordinal) || err.Contains("UNIQUE", StringComparison.OrdinalIgnoreCase) || err.Contains("duplicate", StringComparison.OrdinalIgnoreCase) || err.Contains("2601", StringComparison.Ordinal) || err.Contains("2627", StringComparison.Ordinal)) { return await ExecUpdateAsync(db, p, ct); } throw; } } private static async Task ExistsAsync( ISqlSugarClient db, string table, Dictionary keys, CancellationToken ct) { var where = string.Join(" AND ", keys.Keys.Select((c, i) => $"[{c}]=@k{i}")); var pars = keys.Select((kv, i) => new SugarParameter($"@k{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray(); var sql = $"SELECT TOP 1 1 FROM [{table}] WHERE {where}"; var obj = await db.Ado.GetScalarAsync(sql, pars); return obj != null && obj != DBNull.Value; } /// /// JSON 反序列化常把日期变成带时区的字符串;SQL Server datetime 转换会失败。 /// 可解析则转为 (取本地/Unspecified)。 /// private static object? CoerceDbValue(object? value) { if (value == null) return null; if (value is DateTime or DateTimeOffset or bool or byte or short or int or long or float or double or decimal) return value is DateTimeOffset dto ? dto.LocalDateTime : value; if (value is not string s) return value; if (string.IsNullOrWhiteSpace(s)) return s; if (DateTimeOffset.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dto2)) return dto2.LocalDateTime; if (DateTime.TryParse(s, CultureInfo.InvariantCulture, DateTimeStyles.AssumeLocal, out var dt)) return dt; return s; } private static void ValidateTable(string table) { if (string.IsNullOrWhiteSpace(table) || !ColNameRegex.IsMatch(table)) throw new InvalidOperationException($"非法表名:{table}"); if (!AllowedTables.Contains(table)) throw new InvalidOperationException($"表不在回写白名单:{table}"); } private static void ValidateColumns(Dictionary cols, string part) { foreach (var c in cols.Keys) { if (!ColNameRegex.IsMatch(c)) throw new InvalidOperationException($"非法列名({part}):{c}"); } } private static DbPushPayload ParsePayload(string? json) { using var doc = JsonDocument.Parse(string.IsNullOrWhiteSpace(json) ? "{}" : json); var root = doc.RootElement; var p = new DbPushPayload { Op = root.TryGetProperty("op", out var op) ? op.GetString() ?? "" : "", Table = root.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "", Keys = ReadDict(root, "keys"), Insert = ReadDict(root, "insert"), Update = ReadDict(root, "update"), Expect = ReadDict(root, "expect") }; return p; } private static Dictionary ReadDict(JsonElement root, string name) { var dict = new Dictionary(StringComparer.OrdinalIgnoreCase); if (!root.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Object) return dict; foreach (var prop in el.EnumerateObject()) dict[prop.Name] = JsonElementToObject(prop.Value); return dict; } private static object? JsonElementToObject(JsonElement el) => el.ValueKind switch { JsonValueKind.Null or JsonValueKind.Undefined => null, JsonValueKind.String => el.GetString(), JsonValueKind.True => true, JsonValueKind.False => false, JsonValueKind.Number => el.TryGetInt64(out var l) ? l : el.GetDecimal(), _ => el.GetRawText() }; private sealed class DbPushPayload { public string Op { get; set; } = ""; public string Table { get; set; } = ""; public Dictionary Keys { get; set; } = new(); public Dictionary Insert { get; set; } = new(); public Dictionary Update { get; set; } = new(); public Dictionary Expect { get; set; } = new(); } }