| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152 |
- using Admin.NET.Plugin.AiDOP.DataPlatform;
- using Admin.NET.Plugin.AiDOP.DataPlatform.Wms;
- using Admin.NET.Plugin.AiDOP.DataPlatform.Schema;
- using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
- using Admin.NET.Plugin.AiDOP.Infrastructure;
- using Microsoft.Extensions.Logging;
- using Microsoft.Extensions.Options;
- using SqlSugar;
- namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
- /// <summary>
- /// S5 库存冷链:165 LocationDetail / InvTransHist → stg → std。
- /// 余额事实源仅 LocationDetail;InvMaster 不入标准层。
- /// </summary>
- public sealed class InventoryMdpSyncService : ITransient
- {
- /// <summary>跨实例互斥锁键;真正的取/放锁生命周期见 <see cref="InventoryInboundLockGuard"/>。</summary>
- public const string LockKey = InventoryInboundLockGuard.LockKey;
- private const string LocationEntity = "S5_LOCATION_DETAIL_SQLSERVER";
- private const string TransEntity = "S5_INV_TRANS_HIST_SQLSERVER";
- private const string SourceCodeDefault = "DOPDEMORQ_SQLSERVER";
- private readonly ISqlSugarClient _db;
- private readonly MdpSourcePullDispatcher _pullDispatcher;
- private readonly MdpSourceScopeFactory _scopeFactory;
- private readonly SourceDomainTenantResolver _domainTenant;
- private readonly AidopInventoryOptions _opt;
- private readonly ILogger _logger;
- private readonly MdpNeutralSourceGate _neutralGate;
- private readonly ShipTransNeutralProjection _ship;
- private readonly NativeNeutralProjectionService _native;
- private readonly NeutralRequiredColumnsMonitor _requiredColumns;
- public InventoryMdpSyncService(
- ISqlSugarClient db,
- MdpSourcePullDispatcher pullDispatcher,
- MdpSourceScopeFactory scopeFactory,
- SourceDomainTenantResolver domainTenant,
- IOptions<AidopInventoryOptions> opt,
- ILoggerFactory loggerFactory,
- MdpNeutralSourceGate neutralGate,
- ShipTransNeutralProjection ship,
- NativeNeutralProjectionService native,
- NeutralRequiredColumnsMonitor requiredColumns)
- {
- _db = db;
- _pullDispatcher = pullDispatcher;
- _scopeFactory = scopeFactory;
- _domainTenant = domainTenant;
- _opt = opt.Value;
- _logger = loggerFactory.CreateLogger(nameof(InventoryMdpSyncService));
- _neutralGate = neutralGate;
- _ship = ship;
- _native = native;
- _requiredColumns = requiredColumns;
- }
- public Task<InventorySyncResult> RunBootstrapAsync(CancellationToken cancellationToken = default)
- => RunAsync(bootstrap: true, reconcile: false, cancellationToken);
- public Task<InventorySyncResult> RunIncrementalAsync(CancellationToken cancellationToken = default)
- => RunAsync(bootstrap: false, reconcile: false, cancellationToken);
- public Task<InventorySyncResult> RunReconcileFullAsync(CancellationToken cancellationToken = default)
- => RunAsync(bootstrap: false, reconcile: true, cancellationToken);
- /// <summary>
- /// 仅 stg→std:**不访问源库**。用于链路修复后,复用已完整落地的贴源层让各业务租户
- /// 按当前库位范围口径重新物化标准层。
- /// <para>
- /// 语义为 <b>FULL REPLACE</b>(逐租户、逐正式切片):贴源层是完整历史,
- /// 只跑 UPSERT 既清不掉旧口径残留、也补不齐从未物化过的租户。
- /// 删除范围严格限定「该租户 + 当前正式 source_system + 当前 domain」,
- /// 不触碰 UAT_GENERATOR 等非正式切片。
- /// </para>
- /// </summary>
- public async Task<InventorySyncResult> TransformTransStdFromStgAsync(CancellationToken cancellationToken = default)
- {
- cancellationToken.ThrowIfCancellationRequested();
- var sourceCode = string.IsNullOrWhiteSpace(_opt.SourceCode) ? SourceCodeDefault : _opt.SourceCode.Trim();
- var domain = string.IsNullOrWhiteSpace(_opt.DefaultDomain) ? "8010" : _opt.DefaultDomain.Trim();
- var tenantId = await _domainTenant.ResolveTenantIdAsync(sourceCode, domain, cancellationToken);
- var asOf = DateTime.Now;
- var months = _opt.TransBootstrapMonths <= 0 ? 12 : _opt.TransBootstrapMonths;
- var historyFrom = asOf.Date.AddMonths(-months);
- var batchId = $"S5_INV_XFORM_{asOf:yyyyMMddHHmmss}";
- // await using:异常路径也一定走到释放(DisposeAsync 内先 RELEASE_LOCK 再关专用连接)
- await using var guard = await InventoryInboundLockGuard.TryAcquireAsync(_db, _logger, cancellationToken: cancellationToken);
- if (!guard.Acquired)
- {
- _logger.LogWarning("[InventoryMdpSync] 互斥锁占用,transform-std 本轮跳过 batch={Batch} reason={Reason}",
- batchId, guard.BusyReason);
- return new InventorySyncResult
- {
- BatchId = batchId,
- TenantId = tenantId,
- Domain = domain,
- AsOf = asOf,
- HistoryFrom = historyFrom,
- Skipped = true,
- Message = "lock busy"
- };
- }
- var targetTenants = await ListInventoryScopedTenantsAsync(domain, cancellationToken);
- if (targetTenants.Count == 0)
- {
- // fail closed:没有任何配了合法库位的租户时不写标准层,绝不回落成「按源归属租户写一份」
- _logger.LogWarning(
- "[InventoryMdpSync] domain={Domain} 无任何配置了合法库位的租户,transform-std 本轮不写入", domain);
- return new InventorySyncResult
- {
- BatchId = batchId,
- TenantId = tenantId,
- Domain = domain,
- TransStdRows = 0,
- AsOf = asOf,
- HistoryFrom = historyFrom,
- Skipped = true,
- Message = "no scoped tenant"
- };
- }
- using var xformTimeout = WithLongCommandTimeout();
- var rows = 0;
- foreach (var targetTenantId in targetTenants)
- {
- cancellationToken.ThrowIfCancellationRequested();
- var n = await MdpStdFullReplace.ReplaceAsync(
- _db,
- "mdp_std_inv_trans",
- targetTenantId,
- FormalSliceWhere(sourceCode, domain),
- () => MaterializeInvTransStdAsync(
- tenantId, targetTenantId, batchId: null, asOf, historyFrom, sourceCode),
- cancellationToken);
- rows += n;
- _logger.LogInformation(
- "[InventoryMdpSync] transform-std materialized tenant={Tenant} domain={Domain} rows={Rows}",
- targetTenantId, domain, n);
- }
- _logger.LogInformation(
- "[InventoryMdpSync] transform-std done source={Source} tenants={Cnt} rows={Rows}",
- tenantId, targetTenants.Count, rows);
- return new InventorySyncResult
- {
- BatchId = batchId,
- TenantId = tenantId,
- Domain = domain,
- TransStdRows = rows,
- AsOf = asOf,
- HistoryFrom = historyFrom,
- Message = "OK transform-std"
- };
- }
- /// <summary>
- /// 仅 stg→std(库存余额):**不访问源库**,从一个已完整落地的贴源批次重新物化标准层。
- /// <para>
- /// 用途:同步链路修复后,复用既有完整 stg 批次让各租户按新的库位范围口径重新物化,
- /// 避免为此重新全量拉取源库。
- /// </para>
- /// <para>
- /// 语义为 <b>UPSERT,不做 FULL REPLACE</b>:只按业务键写入/更新本批次覆盖到的行,
- /// 不删除任何既有标准层数据 —— 单个增量批次不代表全量,replace 会造成数据丢失。
- /// 因此本入口<b>不负责</b>清理历史脏快照,那属于全量校准(reconcile)的职责。
- /// </para>
- /// </summary>
- /// <param name="stgBatchId">贴源批次号(sync_batch_id),必须是已完整落地的批次。</param>
- public async Task<InventorySyncResult> TransformInventoryStdFromStgAsync(
- string stgBatchId, CancellationToken cancellationToken = default)
- {
- cancellationToken.ThrowIfCancellationRequested();
- if (string.IsNullOrWhiteSpace(stgBatchId)
- || !System.Text.RegularExpressions.Regex.IsMatch(stgBatchId, @"^[A-Za-z0-9_]+$"))
- throw new InvalidOperationException("非法贴源批次号");
- var sourceCode = string.IsNullOrWhiteSpace(_opt.SourceCode) ? SourceCodeDefault : _opt.SourceCode.Trim();
- var domain = string.IsNullOrWhiteSpace(_opt.DefaultDomain) ? "8010" : _opt.DefaultDomain.Trim();
- var sourceTenantId = await _domainTenant.ResolveTenantIdAsync(sourceCode, domain, cancellationToken);
- var asOf = DateTime.Now;
- // 本入口同样写 mdp_std_inventory,必须与 bootstrap/reconcile/incremental 互斥(原实现漏取锁)
- await using var guard = await InventoryInboundLockGuard.TryAcquireAsync(_db, _logger, cancellationToken: cancellationToken);
- if (!guard.Acquired)
- {
- _logger.LogWarning("[InventoryMdpSync] 互斥锁占用,transform-inventory-std 本轮跳过 batch={Batch} reason={Reason}",
- stgBatchId, guard.BusyReason);
- return new InventorySyncResult
- {
- BatchId = stgBatchId,
- TenantId = sourceTenantId,
- Domain = domain,
- AsOf = asOf,
- Skipped = true,
- Message = "lock busy"
- };
- }
- var stgRows = await _db.Ado.GetIntAsync(
- """
- SELECT COUNT(1) FROM mdp_stg_inventory
- WHERE tenant_id=@TenantId AND source_system=@SourceSystem
- AND source_table='LocationDetail' AND sync_batch_id=@BatchId
- """,
- new List<SugarParameter>
- {
- new("@TenantId", sourceTenantId),
- new("@SourceSystem", sourceCode),
- new("@BatchId", stgBatchId)
- });
- if (stgRows == 0)
- throw new InvalidOperationException(
- $"贴源批次为空或不属于源归属租户:batch={stgBatchId}, sourceTenant={sourceTenantId}");
- var targetTenants = await ListInventoryScopedTenantsAsync(domain, cancellationToken);
- var total = 0;
- foreach (var targetTenantId in targetTenants)
- {
- cancellationToken.ThrowIfCancellationRequested();
- var rows = await InsertInventoryStdAsync(
- targetTenantId, sourceTenantId, stgBatchId, asOf, sourceCode, replaceMode: false);
- total += rows;
- _logger.LogInformation(
- "[InventoryMdpSync] std re-materialized from stg tenant={Tenant} domain={Domain} batch={Batch} rows={Rows}",
- targetTenantId, domain, stgBatchId, rows);
- }
- return new InventorySyncResult
- {
- BatchId = stgBatchId,
- TenantId = sourceTenantId,
- Domain = domain,
- InventoryStdRows = total,
- AsOf = asOf,
- Message = $"OK transform-inventory-std from stg (stgRows={stgRows}, tenants={targetTenants.Count})"
- };
- }
- /// <summary>
- /// 将指定租户已存在的全部库存交易贴源批次转换到标准层。
- /// 供租户级模块重算使用,不拉外部源,也不依赖默认 Domain→Tenant 映射。
- /// <para>
- /// 贴源层归属租户与业务租户不是一回事:正式切片的 stg 挂在
- /// <c>ado_source_domain_tenant_map</c> 解析出的<b>源租户</b>名下,业务归属才按各租户
- /// 库位范围投影(见 <see cref="MaterializeInvTransStdAsync"/>)。因此只按
- /// <c>tenant_id=@TenantId</c> 找贴源会整片漏掉正式切片,该租户的标准层会永远停在
- /// 上一次全量物化的那一刻:之后新增的中立层列(如 <c>approved_time</c>)一直是 NULL,
- /// 依赖它的 KPI 只会算出 NO_DATA 且全程不报错。故这里把正式切片按
- /// 「源租户 stg → 本租户 std」补进来。
- /// </para>
- /// <para>语义为 UPSERT(batchId 传 null 覆盖全部已贴源行),不删任何既有切片。</para>
- /// </summary>
- public async Task<int> TransformTransStdFromStgAsync(
- long tenantId, CancellationToken cancellationToken = default)
- {
- if (tenantId <= 0) throw new ArgumentOutOfRangeException(nameof(tenantId));
- cancellationToken.ThrowIfCancellationRequested();
- var sources = await _db.Ado.SqlQueryAsync<InventoryStageSourceRow>(
- """
- SELECT DISTINCT source_system AS SourceSystem
- FROM mdp_stg_inv_trans
- WHERE tenant_id=@TenantId
- AND source_table='InvTransHist'
- AND NULLIF(TRIM(source_system),'') IS NOT NULL
- """,
- new SugarParameter("@TenantId", tenantId));
- var asOf = DateTime.Now;
- var months = _opt.TransBootstrapMonths <= 0 ? 12 : _opt.TransBootstrapMonths;
- var historyFrom = asOf.Date.AddMonths(-months);
- var slices = sources
- .Where(s => !string.IsNullOrWhiteSpace(s.SourceSystem))
- .Select(s => (SourceTenantId: tenantId, SourceSystem: s.SourceSystem!.Trim()))
- .ToList();
- var sourceCode = string.IsNullOrWhiteSpace(_opt.SourceCode) ? SourceCodeDefault : _opt.SourceCode.Trim();
- var domain = string.IsNullOrWhiteSpace(_opt.DefaultDomain) ? "8010" : _opt.DefaultDomain.Trim();
- var formalSourceTenantId = await _domainTenant.ResolveTenantIdAsync(sourceCode, domain, cancellationToken);
- if (formalSourceTenantId > 0
- && formalSourceTenantId != tenantId
- && !slices.Any(s => string.Equals(s.SourceSystem, sourceCode, StringComparison.OrdinalIgnoreCase)))
- {
- // fail closed:只给「在该 domain 下配了合法库位」的租户补正式切片,
- // 绝不因为解析到了源租户就替无库位范围的租户凭空物化一份。
- var scoped = await ListInventoryScopedTenantsAsync(domain, cancellationToken);
- if (scoped.Contains(tenantId))
- slices.Add((formalSourceTenantId, sourceCode));
- }
- using var xformTimeout = WithLongCommandTimeout();
- var affected = 0;
- foreach (var (sourceTenantId, sourceSystem) in slices)
- {
- cancellationToken.ThrowIfCancellationRequested();
- affected += await MaterializeInvTransStdAsync(
- sourceTenantId, tenantId, batchId: null, asOf, historyFrom, sourceSystem);
- }
- await _requiredColumns.AssertAsync(
- tenantId, "INV_TRANS", "mdp_std_inv_trans", batchId: null, cancellationToken);
- return affected;
- }
- private async Task<InventorySyncResult> RunAsync(bool bootstrap, bool reconcile, CancellationToken cancellationToken)
- {
- cancellationToken.ThrowIfCancellationRequested();
- var sourceCode = string.IsNullOrWhiteSpace(_opt.SourceCode) ? SourceCodeDefault : _opt.SourceCode.Trim();
- var domain = string.IsNullOrWhiteSpace(_opt.DefaultDomain) ? "8010" : _opt.DefaultDomain.Trim();
- // 源归属租户:只决定贴源层(stg)落在谁名下,**不代表业务归属**;
- // 业务归属在标准层物化时按各租户 LocationMaster 库位范围投影决定。
- var tenantId = await _domainTenant.ResolveTenantIdAsync(sourceCode, domain, cancellationToken);
- var asOf = DateTime.Now;
- var batchId = $"S5_INV_{(bootstrap ? "BOOT" : reconcile ? "RECON" : "INCR")}_{asOf:yyyyMMddHHmmss}";
- var months = _opt.TransBootstrapMonths <= 0 ? 12 : _opt.TransBootstrapMonths;
- var historyFrom = asOf.Date.AddMonths(-months);
- // await using:异常路径也一定走到释放(DisposeAsync 内先 RELEASE_LOCK 再关专用连接)
- await using var guard = await InventoryInboundLockGuard.TryAcquireAsync(_db, _logger, cancellationToken: cancellationToken);
- if (!guard.Acquired)
- {
- _logger.LogWarning("[InventoryMdpSync] 跨实例锁占用,本轮跳过 batch={Batch} reason={Reason}", batchId, guard.BusyReason);
- return new InventorySyncResult
- {
- BatchId = batchId,
- TenantId = tenantId,
- Domain = domain,
- Bootstrap = bootstrap,
- AsOf = asOf,
- HistoryFrom = historyFrom,
- Skipped = true,
- Message = "lock busy"
- };
- }
- var upperLoc = await CaptureUpperBoundAsync(sourceCode, "LocationDetail", "UpdateTime", cancellationToken);
- var upperTrans = await CaptureUpperBoundAsync(sourceCode, "InvTransHist", "CreateTime", cancellationToken);
- MdpPullResult locPull;
- if (bootstrap || reconcile)
- {
- var locBatch = $"{batchId}_LOC";
- // NULL UpdateTime 段与非 NULL 段共用同一 batchId,保证 Replace 不漏 NULL 行
- var nullCtx = BuildKeysetCtx(tenantId, locBatch, asOf, historyFrom, upperLoc,
- cursorColumn: "UpdateTime", nullPhase: true, bootstrapFull: true);
- nullCtx.CursorValue = null;
- nullCtx.TieBreakerValue = null;
- await _pullDispatcher.PullAllByEntityCodeAsync(LocationEntity, nullCtx, cancellationToken, maxPages: 50);
- var fullCtx = BuildKeysetCtx(tenantId, locBatch, asOf, historyFrom, upperLoc,
- cursorColumn: "UpdateTime", nullPhase: false, bootstrapFull: true);
- fullCtx.CursorValue = null;
- fullCtx.TieBreakerValue = null;
- fullCtx.BootstrapFrom = null; // LocationDetail 首刷不截时间窗
- locPull = await _pullDispatcher.PullAllByEntityCodeAsync(LocationEntity, fullCtx, cancellationToken, maxPages: 200);
- }
- else
- {
- var incrCtx = BuildKeysetCtx(tenantId, $"{batchId}_LOC", asOf, historyFrom, upperLoc,
- cursorColumn: "UpdateTime", nullPhase: false, bootstrapFull: false);
- // 重叠窗口:从上次游标时间向前回退 LocationOverlapMinutes
- if (!string.IsNullOrWhiteSpace(incrCtx.CursorValue)
- && DateTime.TryParse(incrCtx.CursorValue, out var lastDt))
- {
- var overlap = Math.Max(0, _opt.LocationOverlapMinutes);
- incrCtx.CursorValue = lastDt.AddMinutes(-overlap)
- .ToString("yyyy-MM-dd HH:mm:ss.fff");
- incrCtx.TieBreakerValue = "0";
- }
- locPull = await _pullDispatcher.PullAllByEntityCodeAsync(LocationEntity, incrCtx, cancellationToken, maxPages: 200);
- }
- // —— 标准层按租户物化:一次贴源,逐租户按各自库位范围投影 ——
- // 贴源层是「源+domain」维度(归属 sourceTenantId),标准层是「租户」维度。
- // 拉取游标持久化在 mdp_entity 上、跨租户共享,故绝不能为每个租户各拉一次。
- var targetTenants = await ListInventoryScopedTenantsAsync(domain, cancellationToken);
- if (targetTenants.Count == 0)
- _logger.LogWarning(
- "[InventoryMdpSync] domain={Domain} 无任何配置了合法库位的租户,标准层本轮不写入", domain);
- var inventoryStdRows = 0;
- foreach (var targetTenantId in targetTenants)
- {
- cancellationToken.ThrowIfCancellationRequested();
- int rows;
- if (bootstrap || reconcile)
- {
- rows = await MdpStdFullReplace.ReplaceAsync(
- _db,
- "mdp_std_inventory",
- targetTenantId,
- "source_system='DOPDEMORQ_SQLSERVER'",
- () => InsertInventoryStdAsync(targetTenantId, tenantId, $"{batchId}_LOC", asOf, sourceCode, replaceMode: true),
- cancellationToken);
- }
- else
- {
- rows = await InsertInventoryStdAsync(targetTenantId, tenantId, $"{batchId}_LOC", asOf, sourceCode, replaceMode: false);
- }
- inventoryStdRows += rows;
- _logger.LogInformation(
- "[InventoryMdpSync] std materialized tenant={Tenant} domain={Domain} rows={Rows}",
- targetTenantId, domain, rows);
- }
- var transCtx = BuildKeysetCtx(tenantId, $"{batchId}_TRN", asOf, historyFrom, upperTrans,
- cursorColumn: "CreateTime", nullPhase: false, bootstrapFull: bootstrap || reconcile);
- if (bootstrap || reconcile)
- {
- transCtx.CursorValue = null;
- transCtx.TieBreakerValue = null;
- transCtx.BootstrapFrom = historyFrom;
- }
- var transPull = await _pullDispatcher.PullAllByEntityCodeAsync(TransEntity, transCtx, cancellationToken, maxPages: 500);
- // —— 流水腿与余额腿同构:一次贴源,逐业务租户按各自库位范围投影 ——
- using var transTimeout = WithLongCommandTimeout();
- var transStdRows = 0;
- foreach (var targetTenantId in targetTenants)
- {
- cancellationToken.ThrowIfCancellationRequested();
- int rows;
- if (bootstrap || reconcile)
- {
- rows = await MdpStdFullReplace.ReplaceAsync(
- _db,
- "mdp_std_inv_trans",
- targetTenantId,
- FormalSliceWhere(sourceCode, domain),
- () => MaterializeInvTransStdAsync(
- tenantId, targetTenantId, batchId: null, asOf, historyFrom, sourceCode),
- cancellationToken);
- }
- else
- {
- rows = await MaterializeInvTransStdAsync(
- tenantId, targetTenantId, $"{batchId}_TRN", asOf, historyFrom, sourceCode);
- }
- transStdRows += rows;
- _logger.LogInformation(
- "[InventoryMdpSync] trans std materialized tenant={Tenant} domain={Domain} rows={Rows}",
- targetTenantId, domain, rows);
- }
- _logger.LogInformation(
- "[InventoryMdpSync] done batch={Batch} boot={Boot} recon={Recon} locPulled={LocP} invStd={Inv} trnPulled={TrnP} trnStd={Trn}",
- batchId, bootstrap, reconcile, locPull.RowsPulled, inventoryStdRows, transPull.RowsPulled, transStdRows);
- return new InventorySyncResult
- {
- BatchId = batchId,
- TenantId = tenantId,
- Domain = domain,
- Bootstrap = bootstrap,
- LocationPulled = locPull.RowsPulled,
- LocationWritten = locPull.RowsWritten,
- InventoryStdRows = inventoryStdRows,
- TransPulled = transPull.RowsPulled,
- TransWritten = transPull.RowsWritten,
- TransStdRows = transStdRows,
- AsOf = asOf,
- HistoryFrom = historyFrom,
- Message = "OK"
- };
- }
- private MdpPullContext BuildKeysetCtx(
- long tenantId,
- string batchId,
- DateTime asOf,
- DateTime historyFrom,
- (string? Cursor, string? Tie) upper,
- string cursorColumn,
- bool nullPhase,
- bool bootstrapFull)
- {
- return new MdpPullContext
- {
- TenantId = tenantId,
- BatchId = batchId,
- FullRefresh = false,
- UseKeysetCursor = true,
- CursorColumn = cursorColumn,
- TieBreakerColumn = "RecID",
- UpperCursorValue = upper.Cursor,
- UpperTieBreakerValue = upper.Tie,
- BootstrapFrom = bootstrapFull && cursorColumn == "CreateTime" ? historyFrom : null,
- DeferCursorPersist = false,
- NullTimePhase = nullPhase,
- // 首刷/校准不得继承实体脏水位,否则只会抽到「游标之后」的尾巴
- SkipPersistedKeysetCursor = bootstrapFull || nullPhase,
- StagingOwnerPull = true
- };
- }
- private async Task<(string? Cursor, string? Tie)> CaptureUpperBoundAsync(
- string sourceCode,
- string table,
- string cursorColumn,
- CancellationToken ct)
- {
- if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$")
- || !System.Text.RegularExpressions.Regex.IsMatch(cursorColumn, @"^[A-Za-z0-9_]+$"))
- throw new InvalidOperationException("非法上界查询标识符");
- var remote = await _scopeFactory.GetScopeAsync(sourceCode, ct);
- var rows = await remote.Ado.SqlQueryAsync<UpperRow>(
- $"""
- SELECT TOP 1
- CONVERT(varchar(30), {cursorColumn}, 121) AS CursorText,
- CAST(RecID AS varchar(30)) AS TieText
- FROM {table}
- WHERE {cursorColumn} IS NOT NULL
- ORDER BY {cursorColumn} DESC, RecID DESC
- """);
- var hit = rows.FirstOrDefault();
- return (hit?.CursorText, hit?.TieText);
- }
- /// <summary>
- /// stg → std 物化:**按目标租户的合法库位范围投影**。
- /// <para>
- /// 贴源层(stg)是「源 + domain」维度的全量落地区,归属 <paramref name="sourceTenantId"/>;
- /// 标准层(std)是「租户」维度的可见快照,因此这里必须内联 LocationMaster 做投影:
- /// 只有落在目标租户自己 LocationMaster(同 Domain 且 Typed <> 'Supp')内的库位才写入。
- /// </para>
- /// <para>
- /// 写入不变量:∀ 写入行 → tenant_id = targetTenantId
- /// ∧ location ∈ AllowedLocations(targetTenantId) ∧ domain = 该租户 LocationMaster 的 Domain。
- /// 租户白名单为空 → JOIN 命中 0 行 → 写 0 条(fail closed,绝不退回整个 Domain)。
- /// </para>
- /// <para>
- /// tenant_id 直接取 <paramref name="targetTenantId"/> 而非 MdpJsonSql.TenantFromStg:
- /// 贴源行的 tenant 是「源落地区归属」,不是业务归属,不能顺着传下来。
- /// </para>
- /// </summary>
- private async Task<int> InsertInventoryStdAsync(
- long targetTenantId, long sourceTenantId, string batchId, DateTime asOf, string sourceSystem, bool replaceMode)
- {
- if (!await _neutralGate.AllowsAsync(targetTenantId, "INV_SNAPSHOT", sourceSystem))
- return 0;
- // replaceMode:MdpStdFullReplace 已 DELETE,直接 INSERT;
- // 增量:UPSERT 本批变化。
- var sTenant = "@TargetTenantId";
- var sql = replaceMode
- ? $"""
- INSERT INTO mdp_std_inventory
- (tenant_id, source_system, domain, location, lot_serial, item_num,
- dimension1, dimension2, refs, site, inv_status,
- qty_on_hand, qty_unrestricted, qty_inspection, qty_frozen, qty_available,
- src_rec_id, source_update_time, as_of, sync_batch_id)
- SELECT
- {sTenant},
- @SourceSystem,
- IFNULL({MdpJsonSql.Str("s", "Domain")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Location")}, ''),
- IFNULL({MdpJsonSql.Str("s", "LotSerial")}, ''),
- IFNULL({MdpJsonSql.Str("s", "ItemNum")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Dimension1")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Dimension2")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Refs")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Site")}, ''),
- {MdpJsonSql.Str("s", "InvStatus")},
- IFNULL({MdpJsonSql.Dec("s", "QtyOnHand", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "AvailStatusQty", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "Assay", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "FreezeQty", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "AvailStatusQty", 18, 5)}, 0),
- s.source_row_id,
- {MdpJsonSql.DateTimeSec("s", "UpdateTime")},
- @AsOf,
- @BatchId
- FROM mdp_stg_inventory s
- {TenantScopeJoin}
- WHERE s.tenant_id=@SourceTenantId
- AND s.source_system=@SourceSystem
- AND s.source_table='LocationDetail'
- AND s.sync_batch_id=@BatchId
- """
- : $"""
- INSERT INTO mdp_std_inventory
- (tenant_id, source_system, domain, location, lot_serial, item_num,
- dimension1, dimension2, refs, site, inv_status,
- qty_on_hand, qty_unrestricted, qty_inspection, qty_frozen, qty_available,
- src_rec_id, source_update_time, as_of, sync_batch_id)
- SELECT
- {sTenant},
- @SourceSystem,
- IFNULL({MdpJsonSql.Str("s", "Domain")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Location")}, ''),
- IFNULL({MdpJsonSql.Str("s", "LotSerial")}, ''),
- IFNULL({MdpJsonSql.Str("s", "ItemNum")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Dimension1")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Dimension2")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Refs")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Site")}, ''),
- {MdpJsonSql.Str("s", "InvStatus")},
- IFNULL({MdpJsonSql.Dec("s", "QtyOnHand", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "AvailStatusQty", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "Assay", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "FreezeQty", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "AvailStatusQty", 18, 5)}, 0),
- s.source_row_id,
- {MdpJsonSql.DateTimeSec("s", "UpdateTime")},
- @AsOf,
- @BatchId
- FROM mdp_stg_inventory s
- {TenantScopeJoin}
- WHERE s.tenant_id=@SourceTenantId
- AND s.source_system=@SourceSystem
- AND s.source_table='LocationDetail'
- AND s.sync_batch_id=@BatchId
- ON DUPLICATE KEY UPDATE
- inv_status=VALUES(inv_status),
- qty_on_hand=VALUES(qty_on_hand),
- qty_unrestricted=VALUES(qty_unrestricted),
- qty_inspection=VALUES(qty_inspection),
- qty_frozen=VALUES(qty_frozen),
- qty_available=VALUES(qty_available),
- src_rec_id=VALUES(src_rec_id),
- source_update_time=VALUES(source_update_time),
- as_of=VALUES(as_of),
- sync_batch_id=VALUES(sync_batch_id),
- update_time=CURRENT_TIMESTAMP
- """;
- return await _db.Ado.ExecuteCommandAsync(sql,
- new SugarParameter("@TargetTenantId", targetTenantId),
- new SugarParameter("@SourceTenantId", sourceTenantId),
- new SugarParameter("@SourceSystem", sourceSystem),
- new SugarParameter("@BatchId", batchId),
- new SugarParameter("@AsOf", asOf));
- }
- /// <summary>
- /// 租户库位范围内联投影:贴源行只有落在目标租户自己的合法库位(同 Domain、Typed <> 'Supp')
- /// 才允许进入标准层。这是标准层写入侧的租户安全边界。
- /// <para>
- /// **余额腿与流水腿共用本实现**,只有贴源 JSON 里的库位字段名不同
- /// (LocationDetail 用 <c>Location</c>,InvTransHist 用 <c>Loc</c>)。
- /// 禁止再手写第二套语义略有差异的 LocationMaster JOIN。
- /// </para>
- /// </summary>
- /// <param name="locationJsonKey">贴源 raw_data 中的库位字段名。</param>
- private static string TenantScopeJoinSql(string locationJsonKey) =>
- $"""
- INNER JOIN LocationMaster lm
- ON lm.tenant_id = @TargetTenantId
- AND lm.Domain = IFNULL({MdpJsonSql.Str("s", "Domain")}, '')
- AND lm.location = IFNULL({MdpJsonSql.Str("s", locationJsonKey)}, '')
- AND IFNULL(lm.typed, '') <> 'Supp'
- AND TRIM(lm.location) <> ''
- """;
- /// <summary>库存余额贴源的库位字段名。</summary>
- private static readonly string TenantScopeJoin = TenantScopeJoinSql("Location");
- /// <summary>进出存流水贴源的库位字段名(InvTransHist 用 Loc)。</summary>
- private static readonly string TransTenantScopeJoin = TenantScopeJoinSql("Loc");
- /// <summary>
- /// 库位角色左联:165 的事务码(iss-tr / rct-tr / rct-wo)单看码判不出阶段,
- /// 必须结合库位角色与数量方向,见 <see cref="NeutralTransTypeCodes.LocationRules"/>。
- /// 角色配置缺失时取 UNKNOWN —— 规则一律不成立,流水按「库位角色未配置」隔离,不猜。
- /// </summary>
- private static readonly string LocationRoleJoin =
- $"""
- LEFT JOIN mdp_location_role lr
- ON lr.tenant_id = @TargetTenantId
- AND lr.domain = IFNULL({MdpJsonSql.Str("s", "Domain")}, '')
- AND lr.location = IFNULL({MdpJsonSql.Str("s", "Loc")}, '')
- """;
- /// <summary>库位角色表达式(未配置 → UNKNOWN)。</summary>
- private static readonly string LocationRoleExpr =
- $"IFNULL(lr.location_role, '{NeutralTransTypeCodes.UnknownRole}')";
- /// <summary>
- /// FULL Replace 的删除范围:**只替换「当前租户 + 当前正式源 + 当前 domain」这一份物化投影**。
- /// 绝不能只按 tenant_id 删 —— 那会连带删掉其它来源切片。
- /// UAT 名下的 UAT_GENERATOR 演示行由 1.0.565 清理,不在这次物化删除范围里。
- /// sourceSystem/domain 均来自服务端配置(<c>AidopInventoryOptions</c>),非用户输入;
- /// 仍做标识符白名单校验,杜绝任何拼接注入。
- /// </summary>
- /// <summary>
- /// 标准层物化的命令超时(秒)。默认 30s 不够:正式切片 FULL REPLACE 单次要
- /// DELETE 40 万+ 行再 INSERT ... SELECT 40 万+ 行,实测 DELETE 一步就超时,
- /// 且超时后连接已断、连 ROLLBACK 都会抛 "Connection must be Open" 掩盖真正的超时异常。
- /// 用 <see cref="WithLongCommandTimeout"/> 在物化区间内临时放宽、退出时还原。
- /// </summary>
- private const int MaterializeCommandTimeoutSeconds = 900;
- /// <summary>物化区间内临时放宽命令超时,Dispose 时还原原值(异常路径同样还原)。</summary>
- private sealed class LongCommandTimeoutScope : IDisposable
- {
- private readonly ISqlSugarClient _db;
- private readonly int _original;
- public LongCommandTimeoutScope(ISqlSugarClient db, int seconds)
- {
- _db = db;
- _original = db.Ado.CommandTimeOut;
- db.Ado.CommandTimeOut = seconds;
- }
- public void Dispose() => _db.Ado.CommandTimeOut = _original;
- }
- private LongCommandTimeoutScope WithLongCommandTimeout()
- => new(_db, MaterializeCommandTimeoutSeconds);
- private static string FormalSliceWhere(string sourceSystem, string domain)
- {
- if (!System.Text.RegularExpressions.Regex.IsMatch(sourceSystem, "^[A-Za-z0-9_]+$"))
- throw new InvalidOperationException($"非法 source_system:{sourceSystem}");
- if (!System.Text.RegularExpressions.Regex.IsMatch(domain, "^[A-Za-z0-9_]+$"))
- throw new InvalidOperationException($"非法 domain:{domain}");
- return $"source_system='{sourceSystem}' AND domain='{domain}'";
- }
- /// <summary>
- /// 枚举该 domain 下**拥有合法库存范围**的租户:即在 LocationMaster 里配了非 Supp 库位的启用租户。
- /// 没有库位范围的租户(如默认租户)不会被物化,标准层里不会出现它的快照。
- /// </summary>
- private async Task<List<long>> ListInventoryScopedTenantsAsync(string domain, CancellationToken ct)
- {
- return await _db.Ado.SqlQueryAsync<long>(
- """
- SELECT DISTINCT lm.tenant_id
- FROM LocationMaster lm
- JOIN SysTenant t ON t.Id = lm.tenant_id AND t.Status = 1
- WHERE lm.Domain = @Domain
- AND IFNULL(lm.typed,'') <> 'Supp'
- AND TRIM(lm.location) <> ''
- ORDER BY lm.tenant_id
- """,
- new List<SugarParameter> { new("@Domain", domain) });
- }
- /// <summary>
- /// 进出存流水 stg→std **按业务租户物化**:一次贴源、逐目标租户按各自库位范围投影。
- /// <para>
- /// 修复前本方法只有一个 <c>tenantId</c> 参数,把「源归属租户」当成了「业务租户」:
- /// <c>ado_source_domain_tenant_map</c> 把 DOPDEMORQ_SQLSERVER/8010 登记在 797 名下,
- /// 于是全部流水标准层都写 797,UAT(838257186181189) 名下恒为 0 —— 而正式贴源层里
- /// 落在 UAT 18 个合法库位上的流水实测有 188,419 行。余额腿(InsertInventoryStdAsync)
- /// 早已是「逐租户 + TenantScopeJoin」,本方法此前漏了这一步,本次对齐。
- /// </para>
- /// </summary>
- /// <param name="sourceTenantId">贴源层归属租户(决定读哪批 stg),非业务归属。</param>
- /// <param name="targetTenantId">业务租户(决定 std.tenant_id 与库位投影范围)。</param>
- /// <param name="replaceMode">true=已由 MdpStdFullReplace 删除正式切片,直接 INSERT;false=增量 UPSERT。</param>
- private async Task<int> MaterializeInvTransStdAsync(
- long sourceTenantId, long targetTenantId, string? batchId,
- DateTime asOf, DateTime historyFrom, string sourceSystem)
- {
- if (!await _neutralGate.AllowsAsync(targetTenantId, "INV_TRANS", sourceSystem))
- return 0;
- // 物化 SQL 要 LEFT JOIN 库位角色;未跑到 1.0.564 的环境也不能因缺表而整批失败
- await EnsureLocationRoleTableAsync();
- await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_inv_trans");
- // batchId 为空:转换该源租户下全部已贴源 InvTransHist(中断恢复 / 补物化用)
- var batchPred = string.IsNullOrWhiteSpace(batchId)
- ? "1=1"
- : "s.sync_batch_id=@BatchId";
- var syncBatchExpr = string.IsNullOrWhiteSpace(batchId)
- ? "s.sync_batch_id"
- : "@BatchId";
- // 业务归属 = 目标租户,不再从 stg 反推
- var sTenant = "@TargetTenantId";
- var sql =
- $"""
- INSERT INTO mdp_std_inv_trans
- (tenant_id, source_system, written_by, domain, src_rec_id, trans_type, src_trans_type_raw, biz_doc_type,
- approved_flag, void_flag, summary_flag, approved_time,
- item_num, lot_serial, location,
- dimension1, dimension2, refs, site, qty_change, begin_balance, end_balance,
- eff_date, trans_time, ord_nbr, work_ord, ref_task_no, doc_qty, shipper_num, ship_type, reason, remark, create_user,
- history_from, as_of, sync_batch_id)
- SELECT
- {sTenant},
- @SourceSystem,
- 'DB_SYNC',
- IFNULL({MdpJsonSql.Str("s", "Domain")}, ''),
- s.source_row_id,
- {NeutralTransTypeCodes.TransTypeLookup(
- MdpJsonSql.Str("s", "TransType"),
- LocationRoleExpr,
- $"IFNULL({MdpJsonSql.Dec("s", "QtyChange", 18, 5)}, 0)")},
- {MdpJsonSql.Str("s", "TransType")},
- {NeutralTransTypeCodes.BizDocCase(MdpJsonSql.Str("s", "TransType"), "'OTHER'")},
- 1, 0, 0,
- {MdpJsonSql.DateTimeSec("s", "CreateTime")},
- {MdpJsonSql.Str("s", "ItemNum")},
- {MdpJsonSql.Str("s", "LotSerial")},
- {MdpJsonSql.Str("s", "Loc")},
- IFNULL({MdpJsonSql.Str("s", "Dimension1")}, ''),
- IFNULL({MdpJsonSql.Str("s", "Dimension2")}, ''),
- {MdpJsonSql.Str("s", "Refs")},
- {MdpJsonSql.Str("s", "Site")},
- IFNULL({MdpJsonSql.Dec("s", "QtyChange", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "BeginBalance", 18, 5)}, 0),
- IFNULL({MdpJsonSql.Dec("s", "BeginBalance", 18, 5)}, 0) + IFNULL({MdpJsonSql.Dec("s", "QtyChange", 18, 5)}, 0),
- {MdpJsonSql.DateTimeSec("s", "EffDate")},
- {MdpJsonSql.DateTimeSec("s", "CreateTime")},
- {MdpJsonSql.Str("s", "OrdNbr")},
- {MdpJsonSql.Str("s", "WorkOrd")},
- {MdpJsonSql.Str("s", "WorkOrd")},
- NULLIF({MdpJsonSql.Dec("s", "QtyRequired", 18, 5)}, 0),
- {MdpJsonSql.Str("s", "ShipperNum")},
- {MdpJsonSql.Str("s", "ShipType")},
- {MdpJsonSql.Str("s", "Reason")},
- {MdpJsonSql.Str("s", "Remark")},
- {MdpJsonSql.Str("s", "CreateUser")},
- @HistoryFrom,
- @AsOf,
- {syncBatchExpr}
- FROM mdp_stg_inv_trans s
- {TransTenantScopeJoin}
- {LocationRoleJoin}
- WHERE s.tenant_id=@SourceTenantId
- AND s.source_system=@SourceSystem
- AND s.source_table='InvTransHist'
- AND {batchPred}
- AND {ProjectionGateSql($"IFNULL({MdpJsonSql.Str("s", "Domain")}, '')")}
- ON DUPLICATE KEY UPDATE
- trans_type=VALUES(trans_type),
- src_trans_type_raw=VALUES(src_trans_type_raw),
- biz_doc_type=VALUES(biz_doc_type),
- approved_flag=VALUES(approved_flag),
- approved_time=VALUES(approved_time),
- item_num=VALUES(item_num),
- lot_serial=VALUES(lot_serial),
- location=VALUES(location),
- qty_change=VALUES(qty_change),
- begin_balance=VALUES(begin_balance),
- end_balance=VALUES(end_balance),
- eff_date=VALUES(eff_date),
- trans_time=VALUES(trans_time),
- ord_nbr=VALUES(ord_nbr),
- work_ord=VALUES(work_ord),
- ref_task_no=VALUES(ref_task_no),
- doc_qty=VALUES(doc_qty),
- shipper_num=VALUES(shipper_num),
- ship_type=VALUES(ship_type),
- reason=VALUES(reason),
- remark=VALUES(remark),
- create_user=VALUES(create_user),
- as_of=VALUES(as_of),
- sync_batch_id=VALUES(sync_batch_id),
- update_time=CURRENT_TIMESTAMP
- """;
- var ps = new List<SugarParameter>
- {
- new("@SourceTenantId", sourceTenantId),
- new("@TargetTenantId", targetTenantId),
- new("@SourceSystem", sourceSystem),
- new("@AsOf", asOf),
- new("@HistoryFrom", historyFrom)
- };
- if (!string.IsNullOrWhiteSpace(batchId))
- ps.Add(new SugarParameter("@BatchId", batchId));
- var affected = await _db.Ado.ExecuteCommandAsync(sql, ps);
- await RegisterUnmappedAsync(targetTenantId, sourceSystem, batchId);
- await LogWorkOrderNotFoundAsync(targetTenantId, sourceSystem, batchId);
- if (await _neutralGate.AllowsAsync(targetTenantId, "INV_BAL_MONTHLY", sourceSystem, syncBatchId: batchId))
- await Materialize165InventoryBalanceMonthlyAsync(sourceTenantId, targetTenantId, sourceSystem, batchId);
- if (await _neutralGate.AllowsAsync(targetTenantId, "INV_TRANS", sourceSystem, syncBatchId: batchId))
- {
- await _ship.ProjectAsync(targetTenantId, batchId);
- // 发货流水刚落地,顺手把销售订单行的已交货量与关行状态刷成最新,
- // 否则 S7 订单发货周期要等下一次 S1 重算才看得到本次发货。
- await _native.RefreshSalesLineCompletionAsync(targetTenantId, batchId);
- }
- return affected;
- }
- /// <summary>rct-wo / iss-wo 对不上本租户自建生产任务时记 WO_NOT_FOUND,不改流水。</summary>
- private async Task LogWorkOrderNotFoundAsync(long targetTenantId, string sourceSystem, string? batchId)
- {
- if (!string.Equals(sourceSystem, "DOPDEMORQ_SQLSERVER", StringComparison.OrdinalIgnoreCase))
- return;
- if (!await _neutralGate.AllowsAsync(targetTenantId, "WO_LINE_PROD", MdpSourceIdentity.Native, syncBatchId: batchId))
- return;
- await _db.Ado.ExecuteCommandAsync(
- """
- INSERT INTO mdp_source_gate_log
- (tenant_id, std_object, source_system, gate_reason, row_count, sample_keys, sync_batch_id)
- SELECT @tid, 'WO_LINE_PROD', @src, 'WO_NOT_FOUND', COUNT(*),
- LEFT(GROUP_CONCAT(DISTINCT t.work_ord), 500), @batch
- FROM mdp_std_inv_trans t
- WHERE t.tenant_id=@tid AND t.source_system=@src
- AND t.src_trans_type_raw IN ('rct-wo','iss-wo')
- AND IFNULL(t.work_ord,'')<>''
- AND NOT EXISTS (
- SELECT 1 FROM mdp_std_work_order_line w
- WHERE w.tenant_id=t.tenant_id AND w.source_system='AIDOP_NATIVE'
- AND w.doc_type='PROD_TASK' AND w.order_no=t.work_ord
- )
- HAVING COUNT(*)>0
- """,
- new SugarParameter("@tid", targetTenantId),
- new SugarParameter("@src", sourceSystem),
- new SugarParameter("@batch", batchId ?? ""));
- }
- /// <summary>
- /// 165 月度库存金额:期末数量 × 流水 CurrPrice 的月均,负向数量 × 单价为出库成本。
- /// 无单价的行不进入金额表。
- /// </summary>
- private async Task Materialize165InventoryBalanceMonthlyAsync(
- long sourceTenantId, long targetTenantId, string sourceSystem, string? batchId)
- {
- await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_inventory_balance_monthly");
- var syncBatch = string.IsNullOrWhiteSpace(batchId) ? "INV_BAL_165" : batchId;
- await _db.Ado.ExecuteCommandAsync(
- """
- INSERT INTO mdp_std_inventory_balance_monthly
- (tenant_id, factory_id, source_system, written_by, domain, period_ym,
- category_code, category_name, warehouse_code, warehouse_name, item_code,
- avg_balance_amount, issue_cost_amount, source_biz_key, sync_batch_id, sync_time)
- SELECT
- tenant_id, 1, source_system, 'DB_SYNC', domain, period_ym,
- '', '', warehouse, warehouse, item,
- AVG(end_balance * price),
- SUM(CASE WHEN qty_change < 0 THEN -qty_change * price ELSE 0 END),
- LEFT(CONCAT(domain, ':', period_ym, ':', warehouse, ':', item), 200),
- @batch, NOW()
- FROM (
- SELECT t.tenant_id, t.source_system, IFNULL(t.domain,'') AS domain,
- DATE_FORMAT(t.eff_date, '%Y%m') AS period_ym,
- LEFT(IFNULL(t.location,''), 50) AS warehouse,
- LEFT(IFNULL(t.item_num,''), 100) AS item,
- t.end_balance, t.qty_change,
- CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CurrPrice')),'null') AS DECIMAL(18,6)) AS price
- FROM mdp_std_inv_trans t
- JOIN mdp_stg_inv_trans s
- ON s.tenant_id=@sourceTenant AND s.source_system=t.source_system
- AND s.source_table='InvTransHist' AND s.source_row_id=t.src_rec_id
- WHERE t.tenant_id=@targetTenant AND t.source_system=@source
- AND t.eff_date IS NOT NULL
- ) x
- WHERE price IS NOT NULL AND price <> 0 AND period_ym IS NOT NULL AND period_ym <> ''
- GROUP BY tenant_id, source_system, domain, period_ym, warehouse, item
- ON DUPLICATE KEY UPDATE
- avg_balance_amount=VALUES(avg_balance_amount),
- issue_cost_amount=VALUES(issue_cost_amount),
- written_by=VALUES(written_by),
- sync_batch_id=VALUES(sync_batch_id),
- sync_time=VALUES(sync_time)
- """,
- new
- {
- sourceTenant = sourceTenantId,
- targetTenant = targetTenantId,
- source = sourceSystem,
- batch = syncBatch
- });
- }
- /// <summary>
- /// 库位角色配置表兜底建表。权威定义与初值生成在 1.0.564.sql;
- /// 这里只保证「表存在」,不写初值——角色缺失时阶段码留空并进隔离,不猜。
- /// </summary>
- private async Task EnsureLocationRoleTableAsync() =>
- await MdpSchemaAligner.ExecuteAsync(_db,
- """
- CREATE TABLE IF NOT EXISTS mdp_location_role (
- id BIGINT AUTO_INCREMENT PRIMARY KEY,
- tenant_id BIGINT NOT NULL DEFAULT 0,
- domain VARCHAR(50) NOT NULL DEFAULT '',
- location VARCHAR(100) NOT NULL,
- location_role VARCHAR(24) NOT NULL DEFAULT 'UNKNOWN',
- role_source VARCHAR(24) NOT NULL DEFAULT 'AUTO_DESCR',
- remark VARCHAR(255) NULL,
- create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
- update_time DATETIME NULL ON UPDATE CURRENT_TIMESTAMP,
- UNIQUE KEY uk_location_role (tenant_id, domain, location),
- KEY idx_location_role_role (tenant_id, location_role)
- ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='库位中立角色(阶段码判定用)'
- """);
- /// <summary>
- /// 未能映射到中立阶段码的流水登记隔离表,按原因分级记日志。
- /// <para>
- /// 三类语义不同,禁止一律 Warning:
- /// <list type="bullet">
- /// <item><b>非阶段事件</b>(<see cref="NeutralTransTypeCodes.NonStageCodes"/>,如冻结/解冻):
- /// 不登记、不记日志 —— 它们永远不该出现在待办里。</item>
- /// <item><b>待核验</b>(PENDING_VERIFY):登记,只记 Information。</item>
- /// <item><b>库位角色未配置</b>(LOCATION_ROLE_UNCONFIGURED)/ <b>真未知</b>(UNKNOWN_CODE):
- /// 登记 + Warning,这两类是真有待办动作的。</item>
- /// </list>
- /// </para>
- /// 隔离表尚未建出时只记 Warning,不打断库存物化。
- /// </summary>
- /// <summary>
- /// 业务投影收口。映射表 uk(source_code, domain) 只登记源归属租户,不能再插一行代表 UAT,
- /// 否则会覆盖 797 并让 SourceDomainTenantResolver 报「不唯一」。
- /// 允许投影的目标:映射表上的源归属租户,或该 domain 下自有非 Supp 库位的启用租户。
- /// 默认租户 1300000000001 不是源归属、也没有库位主数据,两条都不成立,历史残留不再增长。
- /// </summary>
- private static string ProjectionGateSql(string domainExpr) =>
- $"""
- (
- EXISTS (
- SELECT 1 FROM ado_source_domain_tenant_map m
- WHERE (m.system_code = @SourceSystem OR m.source_code = @SourceSystem)
- AND m.domain = {domainExpr}
- AND m.tenant_id = @TargetTenantId
- AND m.status = 1
- )
- OR EXISTS (
- SELECT 1 FROM LocationMaster lm
- JOIN SysTenant st ON st.Id = lm.tenant_id AND st.Status = 1
- WHERE lm.tenant_id = @TargetTenantId
- AND IFNULL(lm.Domain,'') = {domainExpr}
- AND IFNULL(lm.typed,'') <> 'Supp'
- AND TRIM(lm.location) <> ''
- )
- )
- """;
- private async Task RegisterUnmappedAsync(long tenantId, string sourceSystem, string? batchId)
- {
- try
- {
- var batchPred = string.IsNullOrWhiteSpace(batchId) ? "1=1" : "t.sync_batch_id=@BatchId";
- var ps = new List<SugarParameter>
- {
- new("@TenantId", tenantId),
- new("@SourceSystem", sourceSystem)
- };
- if (!string.IsNullOrWhiteSpace(batchId))
- ps.Add(new SugarParameter("@BatchId", batchId));
- var pendingScope =
- $"""
- t.trans_type IS NULL AND t.src_trans_type_raw IS NOT NULL
- AND t.src_trans_type_raw NOT IN ({NeutralTransTypeCodes.NonStageInList})
- AND t.src_trans_type_raw NOT IN ({NeutralTransTypeCodes.ExcludedInList})
- """;
- var roleExpr = $"IFNULL(lr.location_role, '{NeutralTransTypeCodes.UnknownRole}')";
- var reasonCase = NeutralTransTypeCodes.UnmappedReasonCase("t.src_trans_type_raw", roleExpr);
- await _db.Ado.ExecuteCommandAsync(
- $"""
- INSERT INTO mdp_std_inv_trans_unmapped
- (tenant_id, source_system, src_rec_id, src_trans_type_raw, unmapped_reason, location, row_count, first_seen, last_seen)
- SELECT t.tenant_id, t.source_system, t.src_rec_id, t.src_trans_type_raw, {reasonCase}, IFNULL(t.location,''), 1, NOW(), NOW()
- FROM mdp_std_inv_trans t
- LEFT JOIN mdp_location_role lr
- ON lr.tenant_id = t.tenant_id
- AND lr.domain = IFNULL(t.domain,'')
- AND lr.location = IFNULL(t.location,'')
- WHERE t.tenant_id=@TenantId AND t.source_system=@SourceSystem
- AND {pendingScope}
- AND {batchPred}
- AND NOT (t.tenant_id = 1300000000001)
- ON DUPLICATE KEY UPDATE last_seen=NOW(), row_count=row_count+1,
- unmapped_reason=VALUES(unmapped_reason), location=VALUES(location)
- """,
- ps);
- var groups = await _db.Ado.SqlQueryAsync<UnmappedTransGroup>(
- $"""
- SELECT t.src_trans_type_raw AS Code, {reasonCase} AS Reason, COUNT(*) AS Cnt
- FROM mdp_std_inv_trans t
- LEFT JOIN mdp_location_role lr
- ON lr.tenant_id = t.tenant_id
- AND lr.domain = IFNULL(t.domain,'')
- AND lr.location = IFNULL(t.location,'')
- WHERE t.tenant_id=@TenantId AND t.source_system=@SourceSystem
- AND {pendingScope}
- AND {batchPred}
- AND NOT (t.tenant_id = 1300000000001)
- GROUP BY t.src_trans_type_raw, {reasonCase}
- """,
- ps);
- foreach (var g in groups)
- {
- if (g.Reason == "PENDING_VERIFY")
- {
- _logger.LogInformation(
- "库存流水有 {Count} 行源事务码 {Code} 语义待核验,已进隔离表暂不计入指标。tenant={TenantId} source={Source}",
- g.Cnt, g.Code, tenantId, sourceSystem);
- continue;
- }
- var action = g.Reason switch
- {
- "LOCATION_ROLE_UNCONFIGURED" => "请在 mdp_location_role 配置该库位角色",
- "OUT_OF_SCOPE_ROLE" => "该库位角色不在此事务码的阶段规则内,不计入指标",
- _ => "请确认该码语义后补入 NeutralTransTypeCodes"
- };
- _logger.LogWarning(
- "库存流水有 {Count} 行源事务码 {Code} 未映射到中立阶段码({Reason}),已进隔离表,指标不计入。{Action}。tenant={TenantId} source={Source}",
- g.Cnt, g.Code, g.Reason, action, tenantId, sourceSystem);
- }
- }
- catch (Exception ex)
- {
- _logger.LogWarning(ex, "未映射流水隔离登记跳过(隔离表未就绪)。tenant={TenantId} source={Source}", tenantId, sourceSystem);
- }
- }
- private sealed class UnmappedTransGroup
- {
- public string Code { get; set; } = "";
- public string Reason { get; set; } = "";
- public int Cnt { get; set; }
- }
- // 取/放锁已迁至 InventoryInboundLockGuard:
- // 原实现在共享 _db(IsAutoCloseConnection=true)上跑 GET_LOCK/RELEASE_LOCK,
- // 命令执行完连接即回池并被驱动 reset,MySQL 当场释放咨询锁 → 跨实例互斥失效。
- private sealed class UpperRow
- {
- public string? CursorText { get; set; }
- public string? TieText { get; set; }
- }
- private sealed class InventoryStageSourceRow
- {
- public string? SourceSystem { get; set; }
- }
- }
- public sealed class InventorySyncResult
- {
- public string BatchId { get; init; } = "";
- public long TenantId { get; init; }
- public string Domain { get; init; } = "";
- public bool Bootstrap { get; init; }
- public int LocationPulled { get; init; }
- public int LocationWritten { get; init; }
- public int InventoryStdRows { get; init; }
- public int TransPulled { get; init; }
- public int TransWritten { get; init; }
- public int TransStdRows { get; init; }
- public DateTime AsOf { get; init; }
- public DateTime? HistoryFrom { get; init; }
- public bool Skipped { get; init; }
- public string? Message { get; init; }
- }
|