using System.Text.Json;
using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
using Admin.NET.Plugin.AiDOP.DataPlatform.HotWatch;
using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace Admin.NET.Plugin.AiDOP.DataPlatform.Wms;
///
/// UAT-INT-001 · Mes 模式:工单物料变更后入队 Outbox,推送时在 165 本地事务内同步一张有效 SM。
/// Local 模式不入队,由调用方继续写本库 Nbr*。
///
public sealed class PickBillNbrSyncService : ITransient
{
public const string ActionCode = "PICK_NBR_SYNC";
public const string TargetSource = CreatePickBillService.TargetSource;
private readonly ISqlSugarClient _db;
private readonly MdpSourceScopeFactory _scopeFactory;
private readonly MdpOutboxEnqueueService _enqueue;
private readonly MdpHotWatchService _hotWatch;
private readonly AidopPickBillOptions _opt;
private readonly ILogger _logger;
public PickBillNbrSyncService(
ISqlSugarClient db,
MdpSourceScopeFactory scopeFactory,
MdpOutboxEnqueueService enqueue,
MdpHotWatchService hotWatch,
IOptions opt,
ILogger logger)
{
_db = db;
_scopeFactory = scopeFactory;
_enqueue = enqueue;
_hotWatch = hotWatch;
_opt = opt.Value;
_logger = logger;
}
public bool IsMes => _opt.IsMes;
///
/// Mes 且工单已下达/投产/暂停时入队(同键刷新 payload)。Local 或未下达返回 false。
///
public async Task TryEnqueueAfterPlanChangeAsync(
long tenantId, string workOrd, string? domain, string account, CancellationToken ct = default)
{
if (!_opt.IsMes)
return false;
workOrd = (workOrd ?? "").Trim();
if (string.IsNullOrWhiteSpace(workOrd))
return false;
var status = await _db.Ado.GetStringAsync(
"""
SELECT IFNULL(LOWER(TRIM(Status)), '') FROM WorkOrdMaster
WHERE tenant_id = @TenantId AND WorkOrd = @WorkOrd
LIMIT 1
""",
new SugarParameter("@TenantId", tenantId),
new SugarParameter("@WorkOrd", workOrd));
if (status is not ("r" or "w" or "s"))
return false;
var mesDomain = ResolveMesDomain(domain);
var idem = $"pick|{mesDomain}|{workOrd}|sm-sync";
if (idem.Length > 200)
idem = idem[..200];
var item = new MdpOutbox
{
TenantId = tenantId,
TargetSourceCode = TargetSource,
ActionCode = ActionCode,
IdemKey = idem,
PayloadJson = JsonSerializer.Serialize(new PickNbrSyncPayload
{
Domain = mesDomain,
WorkOrd = workOrd,
Account = string.IsNullOrWhiteSpace(account) ? "aidop" : account.Trim()
})
};
var ok = await _enqueue.TryEnqueueOrRefreshAsync(item, ct);
_logger.LogInformation(
"[PICK_NBR_SYNC] enqueue tenant={Tenant} wo={Wo} domain={Domain} refreshed={Ok}",
tenantId, workOrd, mesDomain, ok);
return true;
}
public async Task ApplyFromOutboxAsync(MdpOutbox item, CancellationToken ct)
{
PickNbrSyncPayload payload;
try
{
payload = JsonSerializer.Deserialize(item.PayloadJson ?? "")
?? throw new InvalidOperationException("payload 为空");
}
catch (Exception ex)
{
return MdpPushResult.Fail($"PAYLOAD_INVALID: {ex.Message}");
}
var domain = ResolveMesDomain(payload.Domain);
var workOrd = (payload.WorkOrd ?? "").Trim();
var account = Trunc(payload.Account, 24);
if (string.IsNullOrWhiteSpace(workOrd))
return MdpPushResult.Fail("PAYLOAD_INVALID: workOrd 为空");
var plan = await LoadPlanLinesAsync(item.TenantId, workOrd);
ISqlSugarClient ss;
try
{
ss = await _scopeFactory.GetScopeAsync(TargetSource, ct);
}
catch (Exception ex)
{
return MdpPushResult.Fail($"MES_UNREACHABLE: {ex.Message}");
}
try
{
var result = await ApplyOn165Async(ss, item.TenantId, domain, workOrd, account, plan, ct);
return result;
}
catch (Exception ex)
{
_logger.LogWarning(ex, "[PICK_NBR_SYNC] apply failed wo={Wo} domain={Domain}", workOrd, domain);
return MdpPushResult.Fail($"MES_WRITE_FAILED: {ex.Message}");
}
}
private async Task ApplyOn165Async(
ISqlSugarClient ss,
long tenantId,
string domain,
string workOrd,
string account,
List plan,
CancellationToken ct)
{
var sm = (await ss.Ado.SqlQueryAsync(
"""
SELECT TOP 1 RecID, Nbr, Domain
FROM NbrMaster
WHERE Domain = @d AND WorkOrd = @w AND Type = 'SM'
AND ISNULL(TransType, '') = ''
AND ISNULL(IsActive, 1) = 1
ORDER BY RecID DESC
""",
new SugarParameter("@d", domain),
new SugarParameter("@w", workOrd))).FirstOrDefault();
if (sm == null || string.IsNullOrWhiteSpace(sm.Nbr))
{
return MdpPushResult.Skip(JsonSerializer.Serialize(new
{
idempotent = true,
reason = "no_sm",
workOrd,
domain
}));
}
var existing = await ss.Ado.SqlQueryAsync(
"""
SELECT RecID, ItemNum, QtyOrd, QtyRec, CurrQtyOpened, Line, Status
FROM NbrDetail
WHERE NbrRecID = @rid AND Type = 'SM'
AND ISNULL(IsActive, 1) = 1
""",
new SugarParameter("@rid", sm.RecID));
var planMap = plan
.Where(p => !string.IsNullOrWhiteSpace(p.ItemNum))
.ToDictionary(p => p.ItemNum.Trim(), p => p, StringComparer.OrdinalIgnoreCase);
var existingMap = existing
.ToDictionary(x => (x.ItemNum ?? "").Trim(), x => x, StringComparer.OrdinalIgnoreCase);
var now = DateTime.Now;
var updated = 0;
var closed = 0;
var inserted = 0;
await ss.Ado.BeginTranAsync();
try
{
foreach (var line in existing)
{
var key = (line.ItemNum ?? "").Trim();
var still = planMap.TryGetValue(key, out var p);
var newQty = still ? p!.QtyRequired : 0m;
var action = PickBillNbrSyncRules.DecideExisting(line.QtyRec, newQty, still);
if (action == PickBillNbrSyncRules.LineAction.Close)
{
await ss.Ado.ExecuteCommandAsync(
"""
UPDATE NbrDetail
SET Status = 'C',
UpdateUser = @u,
UpdateTime = @t
WHERE RecID = @id
""",
new SugarParameter("@u", account),
new SugarParameter("@t", now),
new SugarParameter("@id", line.RecID));
closed++;
continue;
}
if (action == PickBillNbrSyncRules.LineAction.UpdateQty
&& (line.QtyOrd != newQty || !string.Equals(line.Status ?? "", "", StringComparison.Ordinal)))
{
await ss.Ado.ExecuteCommandAsync(
"""
UPDATE NbrDetail
SET QtyOrd = @q,
CurrQtyOpened = @q,
UM = @um,
ItemName = @name,
Status = '',
UpdateUser = @u,
UpdateTime = @t
WHERE RecID = @id
""",
new SugarParameter("@q", newQty),
new SugarParameter("@um", Trunc(p!.Unit, 8)),
new SugarParameter("@name", Trunc(p.ItemName, 1000)),
new SugarParameter("@u", account),
new SugarParameter("@t", now),
new SugarParameter("@id", line.RecID));
updated++;
}
}
short nextLine = existing.Count > 0 ? (short)(existing.Max(x => x.Line) + 1) : (short)1;
foreach (var p in plan)
{
var key = p.ItemNum.Trim();
if (existingMap.ContainsKey(key))
continue;
await ss.Ado.ExecuteCommandAsync(
"""
INSERT INTO NbrDetail
(Domain, Type, Nbr, Line, ItemNum, Dimension1, Dimension2,
LocationFrom, LocationTo, QtyFrom, QtyTo, UM, [Print], Status, LotSerial,
WorkOrd, QtyOrd, QtyRec, Address, BusinessID, CreateUser, UpdateUser,
CreateTime, UpdateTime, IsActive, IsConfirm, QtyCache, CurrQtyOpened,
IsChanged, NbrRecID, OrdNbr, ItemName, ERPfld1, ERPfld2, OrdLine,
IsGP12Demand, IsGP12Checked, Material, SeqID)
VALUES
(@Domain, 'SM', @Nbr, @Line, @ItemNum, '', '',
@LocationFrom, @LocationTo, 0, 0, @UM, 0, '', @LotSerial,
@WorkOrd, @QtyOrd, 0, '', 0, @User, @User,
@Now, @Now, 1, 0, 0, @QtyOrd,
1, @NbrRecID, '', @ItemName, '', '', 0,
0, 0, 0, 0)
""",
new SugarParameter("@Domain", domain),
new SugarParameter("@Nbr", sm.Nbr),
new SugarParameter("@Line", nextLine++),
new SugarParameter("@ItemNum", Trunc(p.ItemNum, 24)),
new SugarParameter("@LocationFrom", Trunc(p.LocationFrom, 8)),
new SugarParameter("@LocationTo", Trunc(p.LocationTo, 8)),
new SugarParameter("@UM", Trunc(p.Unit, 8)),
new SugarParameter("@LotSerial", p.LotSerial ?? ""),
new SugarParameter("@WorkOrd", workOrd),
new SugarParameter("@QtyOrd", p.QtyRequired),
new SugarParameter("@User", account),
new SugarParameter("@Now", now),
new SugarParameter("@NbrRecID", sm.RecID),
new SugarParameter("@ItemName", Trunc(p.ItemName, 1000)));
inserted++;
}
await ss.Ado.ExecuteCommandAsync(
"""
UPDATE NbrMaster
SET QtyOrd = (
SELECT ISNULL(SUM(QtyOrd), 0)
FROM NbrDetail
WHERE NbrRecID = @rid AND Type = 'SM'
AND ISNULL(IsActive, 1) = 1
AND ISNULL(Status, '') <> 'C'
),
UpdateUser = @u,
UpdateTime = @t
WHERE RecID = @rid
""",
new SugarParameter("@rid", sm.RecID),
new SugarParameter("@u", account),
new SugarParameter("@t", now));
await ss.Ado.CommitTranAsync();
}
catch
{
await ss.Ado.RollbackTranAsync();
throw;
}
try
{
await _hotWatch.EnrollPickBillAsync(domain, sm.Nbr!, workOrd, tenantId, ct);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "[PICK_NBR_SYNC] hot-watch enroll failed nbr={Nbr} wo={Wo}", sm.Nbr, workOrd);
}
var response = JsonSerializer.Serialize(new
{
nbr = sm.Nbr,
workOrd,
domain,
updated,
closed,
inserted
});
_logger.LogInformation(
"[PICK_NBR_SYNC] applied nbr={Nbr} wo={Wo} upd={U} close={C} ins={I}",
sm.Nbr, workOrd, updated, closed, inserted);
return MdpPushResult.Ok(updated + closed + inserted, response);
}
private async Task> LoadPlanLinesAsync(long tenantId, string workOrd)
{
return await _db.Ado.SqlQueryAsync(
"""
SELECT
TRIM(d.ItemNum) AS ItemNum,
SUM(d.QtyRequired) AS QtyRequired,
MAX(IFNULL(d.UM, im.Um)) AS Unit,
MAX(im.Descr) AS ItemName,
MAX(IFNULL(NULLIF(TRIM(d.Location), ''), IFNULL(im.Location, ''))) AS LocationFrom,
MAX(IFNULL(m.Location, '')) AS LocationTo,
MAX(IFNULL(d.LotSerial, '')) AS LotSerial
FROM WorkOrdDetail d
LEFT JOIN WorkOrdMaster m
ON m.WorkOrd = d.WorkOrd AND m.tenant_id = d.tenant_id
LEFT JOIN ItemMaster im
ON im.ItemNum = d.ItemNum AND im.tenant_id = d.tenant_id
WHERE d.tenant_id = @TenantId AND d.WorkOrd = @WorkOrd AND IFNULL(d.IsActive, 0) = 1
GROUP BY TRIM(d.ItemNum)
HAVING SUM(d.QtyRequired) > 0
ORDER BY TRIM(d.ItemNum)
""",
new SugarParameter("@TenantId", tenantId),
new SugarParameter("@WorkOrd", workOrd));
}
/// 165 Domain 最长 8 位;本库若把 tenant_id 写进 Domain 则回落到 8010。
public static string ResolveMesDomain(string? domain)
{
var d = (domain ?? "").Trim();
if (string.IsNullOrWhiteSpace(d) || d.Length > 8)
return "8010";
return d;
}
private static string Trunc(string? s, int max)
{
s ??= "";
return s.Length <= max ? s : s[..max];
}
private sealed class PickNbrSyncPayload
{
public string Domain { get; set; } = "8010";
public string WorkOrd { get; set; } = "";
public string Account { get; set; } = "aidop";
}
private sealed class PlanLine
{
public string ItemNum { get; set; } = "";
public decimal QtyRequired { get; set; }
public string? Unit { get; set; }
public string? ItemName { get; set; }
public string? LocationFrom { get; set; }
public string? LocationTo { get; set; }
public string? LotSerial { get; set; }
}
private sealed class SmHead
{
public int RecID { get; set; }
public string? Nbr { get; set; }
public string? Domain { get; set; }
}
private sealed class SmLine
{
public int RecID { get; set; }
public string? ItemNum { get; set; }
public decimal QtyOrd { get; set; }
public decimal QtyRec { get; set; }
public decimal CurrQtyOpened { get; set; }
public short Line { get; set; }
public string? Status { get; set; }
}
}