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