using System.Collections; using System.Globalization; using System.Text; 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; /// 另支持受限 op=PROC(过程名白名单,见 )。 /// 只保证「不越表、不无条件更新、不注入」;列级白名单由 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", "scm_shd", "scm_shdzb", "scm_shdshph", "MissedPrint", "MissedPrintTransHist", "qms_qcp_insappnentry", "qms_qcp_inspbill", "ScheduleResultOpMaster", "rf_serialnumber", // WP7 基础数据 + 任务(DOP 独占) "LocationMaster", "LocationShelfMaster", "DepartmentMaster", "EmployeeMaster", "EmpWorkDutyMaster", "LinePrinter", "MobileTask", // 单号计数器(NbrSequenceService 直写为主;Outbox 备用) "NbrDayInfo", // S5 IQC QcCheck:resolve 读不良品库位(只读,不写) "PurOrdControl" }; /// /// op=PROC 过程名白名单。本轮仅放 pr_WMS_BPM_SaveInvUpShelf。 /// 依据:doc/plan/S5-IQC合格上架过账(QcCheck)执行任务书.md §3 Q9;白名单变更须同级评审。 /// private static readonly HashSet AllowedProcs = new(StringComparer.OrdinalIgnoreCase) { "pr_WMS_BPM_SaveInvUpShelf" }; 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}"); } var op = (payload.Op ?? "").Trim().ToUpperInvariant(); try { if (op == "PROC") { ValidateProc(payload.Proc); ValidateColumns(payload.Args, "args"); ValidateResolves(payload.Resolves); if (RejectArrayValues(payload.Args, "args") is { } argsErr) return MdpPushResult.Fail(argsErr); } else { ValidateTable(payload.Table); ValidateColumns(payload.Keys, "keys"); ValidateColumns(payload.Insert, "insert"); ValidateColumns(payload.Update, "update"); ValidateColumns(payload.Expect, "expect"); ValidateResolves(payload.Resolves); } var db = await _scopeFactory.GetScopeAsync(source.SourceCode, ct); if (payload.Resolves.Count > 0) { var resolveErr = await ApplyResolvesAsync(db, payload, op, ct); if (resolveErr != null) return MdpPushResult.Fail(resolveErr); } return op switch { "INSERT" => await ExecInsertAsync(db, payload, ct), "UPDATE" => await ExecUpdateAsync(db, payload, item.ActionCode, ct), "UPSERT" => await ExecUpsertAsync(db, payload, item.ActionCode, ct), "PROC" => await ExecProcAsync(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 列"); var arrayErr = RejectArrayValues(p.Keys, "keys") ?? RejectArrayValues(p.Insert, "insert"); if (arrayErr != null) return MdpPushResult.Fail(arrayErr); 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, string? actionCode, CancellationToken ct) { if (p.Keys.Count == 0) throw new InvalidOperationException("UPDATE 禁止无 keys:keys 为空"); if (p.Update.Count == 0) return MdpPushResult.Fail("UPDATE 必须提供 update 列"); var arrayErr = RejectArrayValues(p.Keys, "keys") ?? RejectArrayValues(p.Update, "update"); if (arrayErr != null) return MdpPushResult.Fail(arrayErr); var setParts = p.Update.Keys.Select((c, i) => $"[{c}]=@u{i}").ToList(); var whereParts = p.Keys.Keys.Select((c, i) => $"[{c}]=@k{i}").ToList(); 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)); // P-031 K1-b:expect 标量用 =,数组用 IN;where 与参数共用同一 ei var ei = 0; foreach (var kv in p.Expect) { if (TryAsExpectList(kv.Value, out var list)) { if (list!.Count == 0) return MdpPushResult.Fail("expect 数组为空"); var inNames = new List(list.Count); for (var j = 0; j < list.Count; j++) { var pname = $"@e{ei}_{j}"; inNames.Add(pname); pars.Add(new SugarParameter(pname, CoerceDbValue(list[j]) ?? DBNull.Value)); } whereParts.Add($"[{kv.Key}] IN ({string.Join(",", inNames)})"); ei++; } else { whereParts.Add($"[{kv.Key}]=@e{ei}"); pars.Add(new SugarParameter($"@e{ei}", CoerceDbValue(kv.Value) ?? DBNull.Value)); ei++; } } 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) { // 作废类:目标行已不存在或 expect 不满足(如已投产)→ 幂等成功,避免死信堆积 if (string.Equals(actionCode, "S2_PSD_DEACTIVATE", StringComparison.OrdinalIgnoreCase)) return MdpPushResult.Skip("{\"idempotent\":true,\"reason\":\"no_row_or_expect\"}"); return MdpPushResult.Fail("UPDATE 影响 0 行(行不存在或 expect 不满足)"); } return MdpPushResult.Ok(n); } private static async Task ExecUpsertAsync( ISqlSugarClient db, DbPushPayload p, string? actionCode, 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, actionCode, 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, actionCode, ct); } throw; } } /// /// 调 165 白名单存储过程。关键:部分过程在 @IsProcCall=1 时只开事务不提交, /// 必须由调用方包事务并按 @@TRANCOUNT 收口(见 QcCheck 任务书 §4.2)。 /// private static async Task ExecProcAsync(ISqlSugarClient db, DbPushPayload p, CancellationToken ct) { if (string.IsNullOrWhiteSpace(p.Proc)) return MdpPushResult.Fail("PROC 必须提供 proc"); if (p.Args.Count == 0) return MdpPushResult.Fail("PROC 必须提供 args"); var outputName = string.IsNullOrWhiteSpace(p.Output) ? "ReturnMsg" : p.Output!; if (!ColNameRegex.IsMatch(outputName)) return MdpPushResult.Fail($"非法 OUTPUT 参数名:{outputName}"); // QcCheck:SP 在 @Shelf 仅库位号(@InvShelf='')时,「已是此货架」校验不命中,会重复过账。 // 入队前按标签当前库位预检,已在目标库位则幂等跳过。 var preSkip = await TrySkipQcCheckIfAlreadyThereAsync(db, p, ct); if (preSkip != null) return preSkip; var successNeedle = string.IsNullOrWhiteSpace(p.Success) ? "成功" : p.Success!; var argKeys = p.Args.Keys.ToList(); var assignParts = new List(argKeys.Count + 1); var pars = new List(argKeys.Count); for (var i = 0; i < argKeys.Count; i++) { var key = argKeys[i]; assignParts.Add($"@{key}=@p{i}"); pars.Add(new SugarParameter($"@p{i}", CoerceDbValue(p.Args[key]) ?? DBNull.Value)); } assignParts.Add($"@{outputName}=@ret OUTPUT"); var sql = new StringBuilder(); sql.AppendLine("SET XACT_ABORT ON;"); sql.AppendLine("DECLARE @ret varchar(max) = '';"); sql.AppendLine("BEGIN TRAN;"); sql.Append("EXEC [").Append(p.Proc).Append("] "); sql.Append(string.Join(", ", assignParts)).AppendLine(";"); sql.AppendLine("IF CHARINDEX(N'成功', ISNULL(@ret,'')) = 0"); sql.AppendLine(" WHILE @@TRANCOUNT > 0 ROLLBACK;"); sql.AppendLine("ELSE"); sql.AppendLine(" WHILE @@TRANCOUNT > 0 COMMIT;"); sql.AppendLine("SELECT @ret AS ReturnMsg;"); // 事务收口固定看「成功」(与 SP 文案「保存成功!」对齐);业务成功/幂等再按 payload.success 判定。 var dt = await db.Ado.GetDataTableAsync(sql.ToString(), pars.ToArray()); var retMsg = ""; if (dt is { Rows.Count: > 0 } && dt.Rows[0][0] != DBNull.Value && dt.Rows[0][0] != null) retMsg = Convert.ToString(dt.Rows[0][0]) ?? ""; var response = JsonSerializer.Serialize(new { returnMsg = retMsg }); if (p.Idempotent.Count > 0 && p.Idempotent.Any(s => !string.IsNullOrEmpty(s) && retMsg.Contains(s, StringComparison.Ordinal))) { return MdpPushResult.Skip(response); } if (retMsg.Contains(successNeedle, StringComparison.Ordinal)) return MdpPushResult.Ok(1, response); return MdpPushResult.Fail( string.IsNullOrWhiteSpace(retMsg) ? "PROC 未返回成功消息" : retMsg, response); } /// /// QcCheck 幂等预检:MissedPrint.Location 已等于目标库位时跳过,避免同库位重复 iss-tr-ins。 /// private static async Task TrySkipQcCheckIfAlreadyThereAsync( ISqlSugarClient db, DbPushPayload p, CancellationToken ct) { if (!string.Equals(p.Proc, "pr_WMS_BPM_SaveInvUpShelf", StringComparison.OrdinalIgnoreCase)) return null; if (!p.Args.TryGetValue("TransCode", out var tc) || tc == null) return null; if (!string.Equals(Convert.ToString(tc), "QcCheck", StringComparison.OrdinalIgnoreCase)) return null; if (!p.Args.TryGetValue("Details", out var detailsObj) || detailsObj == null) return null; if (!p.Args.TryGetValue("Shelf", out var shelfObj) || shelfObj == null) return null; var details = Convert.ToString(detailsObj, CultureInfo.InvariantCulture)?.Trim() ?? ""; if (string.IsNullOrEmpty(details) || details.Contains('|', StringComparison.Ordinal) || details.Contains(',', StringComparison.Ordinal)) return null; // 批量细节不做预检,交给 SP var shelf = Convert.ToString(shelfObj, CultureInfo.InvariantCulture)?.Trim() ?? ""; if (string.IsNullOrEmpty(shelf)) return null; var splitIdx = shelf.IndexOf(':'); var locationTo = splitIdx > 0 ? shelf[..splitIdx] : shelf; var curLocObj = await db.Ado.GetScalarAsync( "SELECT TOP 1 [Location] FROM [MissedPrint] WHERE [RecID]=@id", new SugarParameter("@id", details)); var curLoc = curLocObj == null || curLocObj == DBNull.Value ? "" : Convert.ToString(curLocObj, CultureInfo.InvariantCulture)?.Trim() ?? ""; if (string.IsNullOrEmpty(curLoc)) return null; if (string.Equals(curLoc, locationTo, StringComparison.OrdinalIgnoreCase)) { return MdpPushResult.Skip(JsonSerializer.Serialize(new { returnMsg = "已是此货架(执行器预检)", location = curLoc })); } return null; } 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; } /// P-031 K1-b:非 string 的可枚举视为 expect 多值列表。 private static bool TryAsExpectList(object? value, out List? list) { list = null; if (value is null or string) return false; if (value is not IEnumerable enumerable || value is IDictionary) return false; list = new List(); foreach (var item in enumerable) list.Add(item); return true; } /// P-031 K1-c:keys/insert/update 禁止数组值。 private static string? RejectArrayValues(Dictionary cols, string part) { foreach (var kv in cols) { if (kv.Value is null or string) continue; if (kv.Value is IEnumerable and not IDictionary) return $"{part} 不支持数组值,仅 expect 支持"; } return null; } /// /// 目标库自增主键不能由本库 RecID 顶替,外键列须在 165 上现查现填。 /// 上游行还没到位时返回失败,交给 dispatcher 退避重试,等主表推成功后自然收敛。 /// private static async Task ApplyResolvesAsync( ISqlSugarClient db, DbPushPayload p, string op, CancellationToken ct) { foreach (var r in p.Resolves) { var where = string.Join(" AND ", r.Match.Keys.Select((c, i) => $"[{c}]=@r{i}")); var pars = r.Match.Select((kv, i) => new SugarParameter($"@r{i}", CoerceDbValue(kv.Value) ?? DBNull.Value)).ToArray(); var sql = $"SELECT TOP 1 [{r.Column}] FROM [{r.Table}] WHERE {where}"; var val = await db.Ado.GetScalarAsync(sql, pars); if (val == null || val == DBNull.Value) return $"resolve 未命中:{r.Table}.{r.Column}(目标列 {r.TargetColumn}),上游行未就绪"; // PROC:resolve 写入 args;其余写入 insert(既有语义) if (op == "PROC") p.Args[r.TargetColumn] = Convert.ToString(val, CultureInfo.InvariantCulture); else p.Insert[r.TargetColumn] = val; } return null; } private static void ValidateResolves(List resolves) { foreach (var r in resolves) { if (!ColNameRegex.IsMatch(r.TargetColumn)) throw new InvalidOperationException($"非法列名(resolve):{r.TargetColumn}"); ValidateTable(r.Table); if (!ColNameRegex.IsMatch(r.Column)) throw new InvalidOperationException($"非法列名(resolve.column):{r.Column}"); ValidateColumns(r.Match, "resolve.match"); if (r.Match.Count == 0) throw new InvalidOperationException($"resolve.{r.TargetColumn} 缺少 match,禁止无条件取值"); if (RejectArrayValues(r.Match, "resolve.match") is { } err) throw new InvalidOperationException(err); } } 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 ValidateProc(string proc) { if (string.IsNullOrWhiteSpace(proc) || !ColNameRegex.IsMatch(proc)) throw new InvalidOperationException($"非法过程名:{proc}"); if (!AllowedProcs.Contains(proc)) throw new InvalidOperationException($"过程不在回写白名单:{proc}"); } 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() ?? "" : "", Proc = root.TryGetProperty("proc", out var proc) ? proc.GetString() ?? "" : "", Keys = ReadDict(root, "keys"), Insert = ReadDict(root, "insert"), Update = ReadDict(root, "update"), Expect = ReadDict(root, "expect"), Args = ReadDict(root, "args"), Resolves = ReadResolves(root), Output = root.TryGetProperty("output", out var output) ? output.GetString() ?? "" : "", Success = root.TryGetProperty("success", out var success) ? success.GetString() ?? "" : "", Idempotent = ReadStringList(root, "idempotent") }; return p; } private static List ReadStringList(JsonElement root, string name) { var list = new List(); if (!root.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Array) return list; foreach (var item in el.EnumerateArray()) { if (item.ValueKind == JsonValueKind.String) { var s = item.GetString(); if (!string.IsNullOrEmpty(s)) list.Add(s); } } return list; } private static List ReadResolves(JsonElement root) { var list = new List(); if (!root.TryGetProperty("resolve", out var el) || el.ValueKind != JsonValueKind.Object) return list; foreach (var prop in el.EnumerateObject()) { if (prop.Value.ValueKind != JsonValueKind.Object) throw new InvalidOperationException($"resolve.{prop.Name} 必须是对象"); var spec = new ResolveSpec { TargetColumn = prop.Name, Table = prop.Value.TryGetProperty("table", out var t) ? t.GetString() ?? "" : "", Column = prop.Value.TryGetProperty("column", out var c) ? c.GetString() ?? "" : "", Match = ReadDict(prop.Value, "match") }; list.Add(spec); } return list; } 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(), // P-031 K1-a:数组必须保留为列表,否则 expect IN 永远不触发 JsonValueKind.Array => el.EnumerateArray().Select(JsonElementToObject).ToList(), _ => el.GetRawText() }; private sealed class DbPushPayload { public string Op { get; set; } = ""; public string Table { get; set; } = ""; public string Proc { 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(); public Dictionary Args { get; set; } = new(); public List Resolves { get; set; } = new(); public string Output { get; set; } = ""; public string Success { get; set; } = ""; public List Idempotent { get; set; } = new(); } private sealed class ResolveSpec { public string TargetColumn { get; set; } = ""; public string Table { get; set; } = ""; public string Column { get; set; } = ""; public Dictionary Match { get; set; } = new(); } }