Quellcode durchsuchen

feat(mdp): T8 直连 KPI 迁双模式(7迁2留) + S5/S6/S7 统一编排刷新 (server 1.0.266)

- S5/S6/S7 共7个L1 KPI 改读本地标准层 mdp_std_t8_*(经 T8BaseInboundMdpSyncService 源→stg→std),语义等价改写为 MySQL
- S5_L1_002(Cj_Bg_Head_Rep 报表)/S5_L1_004(Rep_总账_存货_V3 TVF) 非干净基表,暂留 legacy 直连 T8
- 方案④:S5MdpRefreshJob 编排 T8 inbound(硬门,失败即中断不降级)→S5→S6→S7 顺序刷新(步级隔离),删除 S6/S7 独立 RefreshJob,保证同 tick 强一致,无新 Job/Cron/锁
- 新增 /t8-base-mdp/inbound 手工补偿端点(partial 文件,未改共享控制器主文件)
- 1.0.266.sql: 幂等启用 6 张 T8 基表 mdp_entity FULL 贴源(biz_key=行级Id)
YY968XX vor 2 Wochen
Ursprung
Commit
bc64119da0

+ 6 - 3
server/Admin.NET.Web.Entry/Admin.NET.Web.Entry.csproj

@@ -11,9 +11,9 @@
     <GenerateSatelliteAssembliesForCore>true</GenerateSatelliteAssembliesForCore>
     <GenerateSatelliteAssembliesForCore>true</GenerateSatelliteAssembliesForCore>
     <Copyright>Admin.NET</Copyright>
     <Copyright>Admin.NET</Copyright>
     <Description>Admin.NET ͨ��Ȩ�޿���ƽ̨</Description>
     <Description>Admin.NET ͨ��Ȩ�޿���ƽ̨</Description>
-    <AssemblyVersion>1.0.265</AssemblyVersion>
-    <FileVersion>1.0.265</FileVersion>
-    <Version>1.0.265</Version>
+    <AssemblyVersion>1.0.266</AssemblyVersion>
+    <FileVersion>1.0.266</FileVersion>
+    <Version>1.0.266</Version>
   </PropertyGroup>
   </PropertyGroup>
 
 
   <ItemGroup>
   <ItemGroup>
@@ -220,6 +220,9 @@
     <None Update="UpdateScripts\1.0.256.sql">
     <None Update="UpdateScripts\1.0.256.sql">
       <CopyToOutputDirectory>Always</CopyToOutputDirectory>
       <CopyToOutputDirectory>Always</CopyToOutputDirectory>
     </None>
     </None>
+    <None Update="UpdateScripts\1.0.266.sql">
+      <CopyToOutputDirectory>Always</CopyToOutputDirectory>
+    </None>
   </ItemGroup>
   </ItemGroup>
 
 
   <ItemGroup>
   <ItemGroup>

+ 31 - 0
server/Admin.NET.Web.Entry/UpdateScripts/1.0.266.sql

@@ -0,0 +1,31 @@
+-- WIP-T8MDP:T8 6 张基表双模式入站启用(提交时 rename 为正式 1.0.<next>.sql 并注册 csproj Copy)
+-- 背景:1.0.178.sql 已登记 mdp_source=T8_V5_SQLSERVER 与 9 个 T8 实体(当时 target_table_name=NULL、sync_mode=NONE,"一期不贴源")。
+-- 本批把 S5_L1_001/003、S6_L1_001/002、S7_L1_001/002/003 这 7 个 KPI 用到的 6 张基表启用为 FULL 贴源,
+--   指向 mdp_stg_t8_*(stg/std 建表由 T8BaseInboundMdpSyncService.EnsureTablesAsync 幂等创建)。
+-- 暂留 legacy 直连的 3 个实体(T8_KC_DD_LIST_CLLIST / T8_CJ_BG_HEAD_REP / T8_REP_ZONGZHANG_CUNHUO_V3)不动,仍 sync_mode=NONE。
+-- 幂等:可重复执行。
+-- 增量策略:FULL 全量。源表极小(12~39 行);MdpDbPullExecutor 的 LastCursor 为字符串字典序比较,
+--   对 bigint Id 增量不安全,且这些表无可靠 lastmodify 列;FULL re-upsert 可带回字段/作废标志更新,且零改通用底座。
+-- source_biz_key:行级 Id(实测 count(*)==count(distinct Id) 唯一,含明细表 kc_tz_list/kc_dd_list/kc_zj_list)。
+
+UPDATE mdp_entity e
+JOIN mdp_source s ON s.id = e.source_id AND s.source_code = 'T8_V5_SQLSERVER'
+SET
+  e.target_table_name = CASE e.entity_code
+    WHEN 'T8_KC_TZ_HEAD' THEN 'mdp_stg_t8_kc_tz_head'
+    WHEN 'T8_KC_TZ_LIST' THEN 'mdp_stg_t8_kc_tz_list'
+    WHEN 'T8_KC_DD_HEAD' THEN 'mdp_stg_t8_kc_dd_head'
+    WHEN 'T8_KC_DD_LIST' THEN 'mdp_stg_t8_kc_dd_list'
+    WHEN 'T8_KC_ZJ_LIST' THEN 'mdp_stg_t8_kc_zj_list'
+    WHEN 'T8_SYS_PELIST' THEN 'mdp_stg_t8_sys_pelist'
+  END,
+  e.sync_mode = 'FULL',
+  e.incr_column = NULL,
+  e.batch_size = 1000,
+  e.biz_key_expr = 'Id',
+  e.status = 1,
+  e.remark = CONCAT(IFNULL(e.remark,''), ' | T8MDP: FULL 贴源启用,biz_key=Id 行级'),
+  e.update_time = NOW()
+WHERE e.entity_code IN (
+  'T8_KC_TZ_HEAD','T8_KC_TZ_LIST','T8_KC_DD_HEAD','T8_KC_DD_LIST','T8_KC_ZJ_LIST','T8_SYS_PELIST'
+);

+ 28 - 0
server/Plugins/Admin.NET.Plugin.AiDOP/Controllers/AidopKanbanController.T8BaseInbound.cs

@@ -0,0 +1,28 @@
+using Admin.NET.Plugin.AiDOP.DataPlatform;
+
+namespace Admin.NET.Plugin.AiDOP.Controllers;
+
+public partial class AidopKanbanController
+{
+    /// <summary>
+    /// T8 基表双模式入站:T8_KC_TZ_HEAD/LIST、T8_KC_DD_HEAD/LIST、T8_KC_ZJ_LIST、T8_SYS_PELIST
+    /// → mdp_stg_t8_* → mdp_std_t8_*。S5/S6/S7 迁移后的 7 个 KPI 读标准层前须先跑本入站。
+    /// 服务经 App.GetRequiredService 解析,避免改动共享控制器构造函数。
+    /// </summary>
+    [HttpPost("t8-base-mdp/inbound")]
+    public async Task<IActionResult> InboundT8BaseMdp(
+        [FromQuery] long? tenantId,
+        [FromQuery] bool fullRefresh = true,
+        [FromQuery] string? entityCode = null,
+        CancellationToken cancellationToken = default)
+    {
+        var svc = App.GetRequiredService<T8BaseInboundMdpSyncService>();
+        var result = await svc.RunInboundAsync(tenantId ?? 0, fullRefresh, entityCode, cancellationToken);
+        return Ok(new
+        {
+            ok = true,
+            result.BatchId,
+            tables = result.Tables.Select(t => new { t.EntityCode, t.StgRows, t.StdRows })
+        });
+    }
+}

+ 285 - 0
server/Plugins/Admin.NET.Plugin.AiDOP/DataPlatform/T8BaseInboundMdpSyncService.cs

@@ -0,0 +1,285 @@
+using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
+using SqlSugar;
+
+namespace Admin.NET.Plugin.AiDOP.DataPlatform;
+
+/// <summary>
+/// T8 基表双模式入站(S5/S6/S7 KPI 共用贴源+标准层)。
+/// 照抄 S6_REPORT 范式:执行器 → mdp_stg_t8_* → transform → mdp_std_t8_*。
+/// 源:mdp_source=T8_V5_SQLSERVER(MdpSourceScopeFactory 复用 Database.json 的 t8_v5 只读连接)。
+/// 实体:mdp_entity T8_KC_TZ_HEAD / T8_KC_TZ_LIST / T8_KC_DD_HEAD / T8_KC_DD_LIST / T8_KC_ZJ_LIST / T8_SYS_PELIST
+///       (种子把 target_table_name 指向对应 mdp_stg_t8_*、sync_mode=FULL、biz_key_expr=Id)。
+/// 只读源、只写 mdp;不写 T8;FULL 全量(源表极小、re-upsert 刷新可带回字段/作废标志更新)。
+/// 下游 S5/S6/S7 KPI 已改读 mdp_std_t8_*;故 KPI 跑批前须先跑一次本入站。
+/// 说明:S5_L1_002(Cj_Bg_Head_Rep)/ S5_L1_004(Rep_总账_存货_V3 TVF)仍走 legacy 直连 T8,
+///       其对应实体 T8_KC_DD_LIST_CLLIST / T8_CJ_BG_HEAD_REP / T8_REP_ZONGZHANG_CUNHUO_V3 本批不贴源。
+/// </summary>
+public sealed class T8BaseInboundMdpSyncService : ITransient
+{
+    private readonly ISqlSugarClient _db;
+    private readonly MdpSourcePullDispatcher _pullDispatcher;
+
+    // 入站实体码(对应 6 张 T8 基表)→ 目标 std 表;stg 表名由 mdp_entity.target_table_name 决定。
+    private static readonly (string EntityCode, string StgTable, string StdTable)[] Entities =
+    {
+        ("T8_KC_TZ_HEAD", "mdp_stg_t8_kc_tz_head", "mdp_std_t8_kc_tz_head"),
+        ("T8_KC_TZ_LIST", "mdp_stg_t8_kc_tz_list", "mdp_std_t8_kc_tz_list"),
+        ("T8_KC_DD_HEAD", "mdp_stg_t8_kc_dd_head", "mdp_std_t8_kc_dd_head"),
+        ("T8_KC_DD_LIST", "mdp_stg_t8_kc_dd_list", "mdp_std_t8_kc_dd_list"),
+        ("T8_KC_ZJ_LIST", "mdp_stg_t8_kc_zj_list", "mdp_std_t8_kc_zj_list"),
+        ("T8_SYS_PELIST", "mdp_stg_t8_sys_pelist", "mdp_std_t8_sys_pelist"),
+    };
+
+    public T8BaseInboundMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher)
+    {
+        _db = db;
+        _pullDispatcher = pullDispatcher;
+    }
+
+    public async Task<T8BaseInboundResult> RunInboundAsync(
+        long tenantId = 0,
+        bool fullRefresh = true,
+        string? entityCode = null,
+        CancellationToken cancellationToken = default)
+    {
+        cancellationToken.ThrowIfCancellationRequested();
+        await EnsureTablesAsync();
+
+        var now = DateTime.Now;
+        var batchId = $"T8_BASE_IN_{now:yyyyMMddHHmmss}";
+        var result = new T8BaseInboundResult { BatchId = batchId };
+
+        var targets = string.IsNullOrWhiteSpace(entityCode)
+            ? Entities
+            : Entities.Where(e => string.Equals(e.EntityCode, entityCode.Trim(), StringComparison.OrdinalIgnoreCase)).ToArray();
+
+        foreach (var (code, stgTable, stdTable) in targets)
+        {
+            cancellationToken.ThrowIfCancellationRequested();
+            var pullCtx = new MdpPullContext
+            {
+                TenantId = tenantId,
+                FullRefresh = fullRefresh,
+                TaskCode = "T8_BASE_MDP_INBOUND",
+                BatchId = $"{batchId}_{code}"
+            };
+            var pull = await _pullDispatcher.PullByEntityCodeAsync(code, pullCtx, cancellationToken);
+            var transformBatch = $"{pullCtx.BatchId}_STD";
+            var stdRows = await TransformStandardAsync(code, stgTable, stdTable, tenantId, transformBatch, now);
+            result.Tables.Add(new T8BaseInboundTableResult
+            {
+                EntityCode = code,
+                StgRows = pull.RowsWritten,
+                StdRows = stdRows
+            });
+        }
+
+        return result;
+    }
+
+    /// <summary>把 PENDING 贴源行按各表 typed 契约投影到 std(JSON key 区分大小写:源标识列为 Id/id,其余小写)。</summary>
+    private async Task<int> TransformStandardAsync(
+        string entityCode, string stgTable, string stdTable, long tenantId, string batchId, DateTime now)
+    {
+        var sql = entityCode switch
+        {
+            "T8_KC_TZ_HEAD" => BuildTzHeadTransform(),
+            "T8_KC_TZ_LIST" => BuildTzListTransform(),
+            "T8_KC_DD_HEAD" => BuildDdHeadTransform(),
+            "T8_KC_DD_LIST" => BuildDdListTransform(),
+            "T8_KC_ZJ_LIST" => BuildZjListTransform(),
+            "T8_SYS_PELIST" => BuildPelistTransform(),
+            _ => throw new InvalidOperationException($"未支持的 T8 入站实体:{entityCode}")
+        };
+
+        var affected = await _db.Ado.ExecuteCommandAsync(sql,
+            new SugarParameter("@tid", tenantId),
+            new SugarParameter("@batch", batchId),
+            new SugarParameter("@now", now));
+
+        await _db.Ado.ExecuteCommandAsync(
+            $"UPDATE {stgTable} SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING'");
+        return affected;
+    }
+
+    // JSON 取值片段:数字/整数经 NULLIF 去 'null';datetime 把 ISO 'T' 换空格再 STR_TO_DATE(尾部小数位被忽略)。
+    private static string JStr(string k) => $"NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.{k}')),'null')";
+    private static string JNum(string k) => $"CAST({JStr(k)} AS DECIMAL(18,6))";
+    private static string JInt(string k) => $"CAST({JStr(k)} AS SIGNED)";
+    // ISO 时间戳(System.Text.Json 输出如 2026-05-08T13:45:16.92):'T'→空格、去亚秒再 STR_TO_DATE,
+    // 否则 MySQL 严格模式对残留 '.92' 报 Truncated incorrect datetime value。KPI 只用日级,丢亚秒无影响。
+    private static string JDate(string k) => $"STR_TO_DATE(SUBSTRING_INDEX(REPLACE({JStr(k)},'T',' '),'.',1),'%Y-%m-%d %H:%i:%s')";
+    private static string SrcId => $"CAST(COALESCE({JStr("Id")},{JStr("id")}) AS SIGNED)";
+
+    private static string StgHead(string table) => $@"
+FROM {table} s
+WHERE s.process_status='PENDING' AND s.source_biz_key IS NOT NULL AND s.source_biz_key<>''";
+
+    private static string BuildTzHeadTransform() => $@"
+INSERT INTO mdp_std_t8_kc_tz_head
+  (tenant_id, source_system, src_id, ztid, lbs, hzyn, zfyn, shyn, shtime, date0,
+   source_row_id, source_biz_key, sync_batch_id, sync_time)
+SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'),
+  {SrcId}, {JStr("ztid")}, {JStr("lbs")}, {JInt("hzyn")}, {JInt("zfyn")}, {JInt("shyn")},
+  {JDate("shtime")}, {JDate("date0")},
+  IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
+{StgHead("mdp_stg_t8_kc_tz_head")}
+ON DUPLICATE KEY UPDATE
+  src_id=VALUES(src_id), ztid=VALUES(ztid), lbs=VALUES(lbs),
+  hzyn=VALUES(hzyn), zfyn=VALUES(zfyn), shyn=VALUES(shyn),
+  shtime=VALUES(shtime), date0=VALUES(date0),
+  source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
+  sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
+
+    private static string BuildTzListTransform() => $@"
+INSERT INTO mdp_std_t8_kc_tz_list
+  (tenant_id, source_system, src_id, idid, code, lynoid, slzx, sl, gdyn, gdtime, rwnoid, jhdate, addtime,
+   source_row_id, source_biz_key, sync_batch_id, sync_time)
+SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'),
+  {SrcId}, {JInt("idid")}, {JStr("code")}, {JStr("lynoid")}, {JNum("slzx")}, {JNum("sl")},
+  {JInt("gdyn")}, {JDate("gdtime")}, {JStr("rwnoid")}, {JDate("jhdate")}, {JDate("addtime")},
+  IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
+{StgHead("mdp_stg_t8_kc_tz_list")}
+ON DUPLICATE KEY UPDATE
+  src_id=VALUES(src_id), idid=VALUES(idid), code=VALUES(code), lynoid=VALUES(lynoid),
+  slzx=VALUES(slzx), sl=VALUES(sl), gdyn=VALUES(gdyn), gdtime=VALUES(gdtime),
+  rwnoid=VALUES(rwnoid), jhdate=VALUES(jhdate), addtime=VALUES(addtime),
+  source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
+  sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
+
+    private static string BuildDdHeadTransform() => $@"
+INSERT INTO mdp_std_t8_kc_dd_head
+  (tenant_id, source_system, src_id, ztid, lbs, zf, shyn, noid,
+   source_row_id, source_biz_key, sync_batch_id, sync_time)
+SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'),
+  {SrcId}, {JStr("ztid")}, {JStr("lbs")}, {JInt("zf")}, {JInt("shyn")}, {JStr("noid")},
+  IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
+{StgHead("mdp_stg_t8_kc_dd_head")}
+ON DUPLICATE KEY UPDATE
+  src_id=VALUES(src_id), ztid=VALUES(ztid), lbs=VALUES(lbs),
+  zf=VALUES(zf), shyn=VALUES(shyn), noid=VALUES(noid),
+  source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
+  sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
+
+    private static string BuildDdListTransform() => $@"
+INSERT INTO mdp_std_t8_kc_dd_list
+  (tenant_id, source_system, src_id, idid, rwnoid, code, sl, slzx, jhdate, gdyn, gdtime, addtime,
+   source_row_id, source_biz_key, sync_batch_id, sync_time)
+SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'),
+  {SrcId}, {JInt("idid")}, {JStr("rwnoid")}, {JStr("code")}, {JNum("sl")}, {JNum("slzx")},
+  {JDate("jhdate")}, {JInt("gdyn")}, {JDate("gdtime")}, {JDate("addtime")},
+  IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
+{StgHead("mdp_stg_t8_kc_dd_list")}
+ON DUPLICATE KEY UPDATE
+  src_id=VALUES(src_id), idid=VALUES(idid), rwnoid=VALUES(rwnoid), code=VALUES(code),
+  sl=VALUES(sl), slzx=VALUES(slzx), jhdate=VALUES(jhdate), gdyn=VALUES(gdyn), gdtime=VALUES(gdtime),
+  addtime=VALUES(addtime),
+  source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
+  sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
+
+    private static string BuildZjListTransform() => $@"
+INSERT INTO mdp_std_t8_kc_zj_list
+  (tenant_id, source_system, src_id, ztid, lyid, zjyn, shdate,
+   source_row_id, source_biz_key, sync_batch_id, sync_time)
+SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'),
+  {SrcId}, {JStr("ztid")}, {JInt("lyid")}, {JInt("zjyn")}, {JDate("shdate")},
+  IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
+{StgHead("mdp_stg_t8_kc_zj_list")}
+ON DUPLICATE KEY UPDATE
+  src_id=VALUES(src_id), ztid=VALUES(ztid), lyid=VALUES(lyid), zjyn=VALUES(zjyn), shdate=VALUES(shdate),
+  source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
+  sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
+
+    private static string BuildPelistTransform() => $@"
+INSERT INTO mdp_std_t8_sys_pelist
+  (tenant_id, source_system, src_id, ztid, zzzt, gw,
+   source_row_id, source_biz_key, sync_batch_id, sync_time)
+SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'),
+  {SrcId}, {JStr("ztid")}, {JStr("zzzt")}, {JStr("gw")},
+  IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
+{StgHead("mdp_stg_t8_sys_pelist")}
+ON DUPLICATE KEY UPDATE
+  src_id=VALUES(src_id), ztid=VALUES(ztid), zzzt=VALUES(zzzt), gw=VALUES(gw),
+  source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
+  sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
+
+    /// <summary>建 6 stg(通用信封契约)+ 6 std(KPI 实际引用的 typed 列)。幂等 IF NOT EXISTS。</summary>
+    private async Task EnsureTablesAsync()
+    {
+        foreach (var (_, stgTable, _) in Entities)
+            await _db.Ado.ExecuteCommandAsync(StgDdl(stgTable));
+
+        await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_tz_head",
+            "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lbs varchar(100) DEFAULT NULL, " +
+            "hzyn int DEFAULT NULL, zfyn int DEFAULT NULL, shyn int DEFAULT NULL, " +
+            "shtime datetime DEFAULT NULL, date0 date DEFAULT NULL"));
+        await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_tz_list",
+            "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL, code varchar(400) DEFAULT NULL, " +
+            "lynoid varchar(400) DEFAULT NULL, slzx decimal(18,6) DEFAULT NULL, sl decimal(18,6) DEFAULT NULL, " +
+            "gdyn bigint DEFAULT NULL, gdtime datetime DEFAULT NULL, rwnoid varchar(400) DEFAULT NULL, " +
+            "jhdate date DEFAULT NULL, addtime datetime DEFAULT NULL"));
+        await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_dd_head",
+            "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lbs varchar(100) DEFAULT NULL, " +
+            "zf int DEFAULT NULL, shyn int DEFAULT NULL, noid varchar(400) DEFAULT NULL"));
+        await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_dd_list",
+            "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL, rwnoid varchar(400) DEFAULT NULL, " +
+            "code varchar(400) DEFAULT NULL, sl decimal(18,6) DEFAULT NULL, slzx decimal(18,6) DEFAULT NULL, " +
+            "jhdate date DEFAULT NULL, gdyn int DEFAULT NULL, gdtime datetime DEFAULT NULL, addtime datetime DEFAULT NULL"));
+        await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_zj_list",
+            "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lyid bigint DEFAULT NULL, " +
+            "zjyn int DEFAULT NULL, shdate datetime DEFAULT NULL"));
+        await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_sys_pelist",
+            "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, zzzt varchar(100) DEFAULT NULL, " +
+            "gw varchar(400) DEFAULT NULL"));
+    }
+
+    private static string StgDdl(string table) => $@"
+CREATE TABLE IF NOT EXISTS {table} (
+  id bigint NOT NULL AUTO_INCREMENT,
+  tenant_id bigint NOT NULL DEFAULT 0,
+  source_system varchar(50) DEFAULT NULL,
+  source_table varchar(200) DEFAULT NULL,
+  source_row_id varchar(200) DEFAULT NULL,
+  source_biz_key varchar(300) DEFAULT NULL,
+  raw_data json DEFAULT NULL,
+  sync_batch_id varchar(100) DEFAULT NULL,
+  create_time datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
+  sync_time datetime DEFAULT CURRENT_TIMESTAMP,
+  process_status varchar(20) NOT NULL DEFAULT 'PENDING',
+  process_message varchar(500) DEFAULT NULL,
+  update_time datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+  PRIMARY KEY (id),
+  UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
+  KEY idx_batch (sync_batch_id)
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='T8 基表贴源层(双模式入站)'";
+
+    private static string StdDdl(string table, string bizCols) => $@"
+CREATE TABLE IF NOT EXISTS {table} (
+  id bigint NOT NULL AUTO_INCREMENT,
+  tenant_id bigint NOT NULL DEFAULT 0,
+  source_system varchar(50) NOT NULL DEFAULT 'T8',
+  {bizCols},
+  source_row_id varchar(200) NOT NULL,
+  source_biz_key varchar(300) NOT NULL,
+  sync_batch_id varchar(100) NOT NULL,
+  sync_time datetime NOT NULL,
+  update_time datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+  PRIMARY KEY (id),
+  UNIQUE KEY uk_{table} (tenant_id, source_system, source_biz_key),
+  KEY idx_{table}_batch (sync_batch_id)
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='T8 基表标准层(双模式入站)'";
+}
+
+/// <summary>T8 基表入站结果。</summary>
+public sealed class T8BaseInboundResult
+{
+    public string BatchId { get; set; } = string.Empty;
+    public List<T8BaseInboundTableResult> Tables { get; } = new();
+}
+
+public sealed class T8BaseInboundTableResult
+{
+    public string EntityCode { get; set; } = string.Empty;
+    public int StgRows { get; set; }
+    public int StdRows { get; set; }
+}

+ 27 - 31
server/Plugins/Admin.NET.Plugin.AiDOP/FinishedWarehouse/S7MdpSyncTransformService.cs

@@ -6,8 +6,9 @@ using System.Text.Json;
 namespace Admin.NET.Plugin.AiDOP.FinishedWarehouse;
 namespace Admin.NET.Plugin.AiDOP.FinishedWarehouse;
 
 
 /// <summary>
 /// <summary>
-/// S7 成品仓储 — T8 KPI 数据底座与刷新转换服务。
-/// 路径 A:方老师 v5.4 KPI 字段对照表 J 列 SQL 原逻辑直发 T8 SQL Server(ConfigId=t8_v5)。
+/// S7 成品仓储 — KPI 计算与刷新转换服务。
+/// 双模式:读本地标准层 mdp_std_t8_kc_*(由 T8BaseInboundMdpSyncService 从 T8 贴源→标准),不再直连 T8。
+/// 计算逻辑沿用方老师 v5.4 KPI J 列口径(等价改写为 MySQL 读 std)。
 /// 包含 KPI:S7_L1_001 订单发货周期 / S7_L1_002 订单发货满足率 / S7_L1_003 成品仓储人效。
 /// 包含 KPI:S7_L1_001 订单发货周期 / S7_L1_002 订单发货满足率 / S7_L1_003 成品仓储人效。
 /// </summary>
 /// </summary>
 public class S7MdpSyncTransformService : ITransient
 public class S7MdpSyncTransformService : ITransient
@@ -17,7 +18,6 @@ public class S7MdpSyncTransformService : ITransient
     private readonly ILogger<S7MdpSyncTransformService> _logger;
     private readonly ILogger<S7MdpSyncTransformService> _logger;
     private const string JobCode = "S7_MDP_SYNC_TRANSFORM";
     private const string JobCode = "S7_MDP_SYNC_TRANSFORM";
     private const string JobName = "S7 成品仓储 MDP 同步与转换";
     private const string JobName = "S7 成品仓储 MDP 同步与转换";
-    private const string T8ConfigId = "t8_v5";
     private const string ModuleCode = "S7";
     private const string ModuleCode = "S7";
     // FAILURE-NOTIFICATION-1:超级管理员 superAdmin.NET(AccountType=999)
     // FAILURE-NOTIFICATION-1:超级管理员 superAdmin.NET(AccountType=999)
     private const long NoticeReceiverUserId = 1300000000101L;
     private const long NoticeReceiverUserId = 1300000000101L;
@@ -93,31 +93,32 @@ public class S7MdpSyncTransformService : ITransient
     {
     {
         var sub = new KpiBuildSubResult();
         var sub = new KpiBuildSubResult();
 
 
+        // 双模式:读本地标准层 mdp_std_t8_*(源 identity Id→src_id);datediff(day,a,b)→DATEDIFF(b,a),IsNull→IFNULL。
         const string sql = @"
         const string sql = @"
 select noid as noid,
 select noid as noid,
-       datediff(day, min(shdate), max(shtime)) as scts
+       datediff(max(shtime), min(shdate)) as scts
  from (
  from (
    select a.noid as noid, b.code as code,
    select a.noid as noid, b.code as code,
-          IsNull(c.shdate, b.addtime) as shdate,
+          ifnull(c.shdate, b.addtime) as shdate,
           (case when b.gdyn=1 then b.gdtime else d.shtime end) as shtime,
           (case when b.gdyn=1 then b.gdtime else d.shtime end) as shtime,
           (case when b.gdyn=1 or b.sl<=b.slzx then 1 else 0 end) as wczt
           (case when b.gdyn=1 or b.sl<=b.slzx then 1 else 0 end) as wczt
-    from kc_dd_head a with(nolock)
-    left join kc_dd_list b with(nolock) on a.Id=b.idid
+    from mdp_std_t8_kc_dd_head a
+    left join mdp_std_t8_kc_dd_list b on a.src_id=b.idid
     left join (
     left join (
-      select min(b.id) as id, b.lynoid as lynoid, b.code as code,
+      select min(b.src_id) as id, b.lynoid as lynoid, b.code as code,
              max(a.shtime) as shtime, sum(b.slzx) as slzx
              max(a.shtime) as shtime, sum(b.slzx) as slzx
-       from kc_tz_head a with(nolock)
-       inner join kc_tz_list b with(nolock) on a.Id=b.idid
+       from mdp_std_t8_kc_tz_head a
+       inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
        where a.ztid=@ztid and a.lbs='销售出库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
        where a.ztid=@ztid and a.lbs='销售出库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
        group by b.lynoid, b.code
        group by b.lynoid, b.code
     ) d on b.rwnoid=d.lynoid and b.code=d.code
     ) d on b.rwnoid=d.lynoid and b.code=d.code
-    left join kc_zj_list c on a.ztid=c.ztid and c.lyid=b.id and c.zjyn=1
+    left join mdp_std_t8_kc_zj_list c on a.ztid=c.ztid and c.lyid=b.src_id and c.zjyn=1
     where a.ztid=@ztid and a.lbs='销售订单' and a.zf=0 and a.shyn=1
     where a.ztid=@ztid and a.lbs='销售订单' and a.zf=0 and a.shyn=1
  ) n
  ) n
  group by noid
  group by noid
  having min(wczt)=1";
  having min(wczt)=1";
 
 
-        var rows = await QueryT8Async<S7CycleRow>(sql, new[] { new SugarParameter("@ztid", option.SourceZtid) });
+        var rows = await _db.Ado.SqlQueryAsync<S7CycleRow>(sql, new[] { new SugarParameter("@ztid", option.SourceZtid) });
         sub.T8Rows = rows.Count;
         sub.T8Rows = rows.Count;
 
 
         var dwdAffected = 0;
         var dwdAffected = 0;
@@ -161,6 +162,7 @@ ON DUPLICATE KEY UPDATE
     {
     {
         var sub = new KpiBuildSubResult();
         var sub = new KpiBuildSubResult();
 
 
+        // 双模式:读标准层;convert(varchar(10),shtime,23)→date(shtime)。
         const string sql = @"
         const string sql = @"
 select noid as noid,
 select noid as noid,
        count(noid) as total_rows,
        count(noid) as total_rows,
@@ -168,13 +170,13 @@ select noid as noid,
  from (
  from (
    select a.noid as noid, b.rwnoid as rwnoid, b.code as code,
    select a.noid as noid, b.rwnoid as rwnoid, b.code as code,
           (case when sum(d.slzx)>=b.sl then 1 else 0 end) as wczt
           (case when sum(d.slzx)>=b.sl then 1 else 0 end) as wczt
-    from kc_dd_head a with(nolock)
-    left join kc_dd_list b with(nolock) on a.Id=b.idid
+    from mdp_std_t8_kc_dd_head a
+    left join mdp_std_t8_kc_dd_list b on a.src_id=b.idid
     left join (
     left join (
       select b.lynoid as lynoid, b.code as code,
       select b.lynoid as lynoid, b.code as code,
-             convert(varchar(10), a.shtime, 23) as shtime, b.slzx as slzx
-       from kc_tz_head a with(nolock)
-       inner join kc_tz_list b with(nolock) on a.Id=b.idid
+             date(a.shtime) as shtime, b.slzx as slzx
+       from mdp_std_t8_kc_tz_head a
+       inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
        where a.ztid=@ztid and a.lbs='销售出库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
        where a.ztid=@ztid and a.lbs='销售出库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
     ) d on b.rwnoid=d.lynoid and b.code=d.code and d.shtime<=b.jhdate
     ) d on b.rwnoid=d.lynoid and b.code=d.code and d.shtime<=b.jhdate
     where a.ztid=@ztid and a.lbs='销售订单' and a.zf=0 and a.shyn=1
     where a.ztid=@ztid and a.lbs='销售订单' and a.zf=0 and a.shyn=1
@@ -182,7 +184,7 @@ select noid as noid,
  ) n
  ) n
  group by noid";
  group by noid";
 
 
-        var rows = await QueryT8Async<S7FulfillmentRow>(sql, new[] { new SugarParameter("@ztid", option.SourceZtid) });
+        var rows = await _db.Ado.SqlQueryAsync<S7FulfillmentRow>(sql, new[] { new SugarParameter("@ztid", option.SourceZtid) });
         sub.T8Rows = rows.Count;
         sub.T8Rows = rows.Count;
 
 
         var dwdAffected = 0;
         var dwdAffected = 0;
@@ -236,14 +238,14 @@ ON DUPLICATE KEY UPDATE
 
 
         const string sqlNumer = @"
         const string sqlNumer = @"
 select b.lynoid as lynoid, b.code as code,
 select b.lynoid as lynoid, b.code as code,
-       convert(varchar(10), a.shtime, 23) as shtime, b.slzx as slzx
- from kc_tz_head a with(nolock)
- inner join kc_tz_list b with(nolock) on a.Id=b.idid
+       date(a.shtime) as shtime, b.slzx as slzx
+ from mdp_std_t8_kc_tz_head a
+ inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
  where a.ztid=@ztid and a.lbs='销售出库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  where a.ztid=@ztid and a.lbs='销售出库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
-  and convert(varchar(10), a.shtime, 23) between @startDateText and @endDateText";
+  and date(a.shtime) between @startDateText and @endDateText";
         const string sqlDenom = @"
         const string sqlDenom = @"
 select count(*) as penum
 select count(*) as penum
- from sys_pelist with(nolock)
+ from mdp_std_t8_sys_pelist
  where ztid=@ztid and zzzt='在职' and gw='仓管'";
  where ztid=@ztid and zzzt='在职' and gw='仓管'";
 
 
         var pNumer = new[]
         var pNumer = new[]
@@ -254,8 +256,8 @@ select count(*) as penum
         };
         };
         var pDenom = new[] { new SugarParameter("@ztid", option.SourceZtid) };
         var pDenom = new[] { new SugarParameter("@ztid", option.SourceZtid) };
 
 
-        var numerRows = await QueryT8Async<S7ShipmentDetailRow>(sqlNumer, pNumer);
-        var denomRows = await QueryT8Async<S7PeNumRow>(sqlDenom, pDenom);
+        var numerRows = await _db.Ado.SqlQueryAsync<S7ShipmentDetailRow>(sqlNumer, pNumer);
+        var denomRows = await _db.Ado.SqlQueryAsync<S7PeNumRow>(sqlDenom, pDenom);
         sub.T8Rows = numerRows.Count + denomRows.Count;
         sub.T8Rows = numerRows.Count + denomRows.Count;
 
 
         decimal? shipmentQty = numerRows.Sum(r => r.slzx ?? 0m);
         decimal? shipmentQty = numerRows.Sum(r => r.slzx ?? 0m);
@@ -308,12 +310,6 @@ ON DUPLICATE KEY UPDATE
 
 
     // ─────────────────────────────────────────────────────────────────────────
     // ─────────────────────────────────────────────────────────────────────────
 
 
-    private async Task<List<T>> QueryT8Async<T>(string sql, SugarParameter[] parameters)
-    {
-        var t8 = _db.AsTenant().GetConnectionScope(T8ConfigId);
-        return await t8.Ado.SqlQueryAsync<T>(sql, parameters);
-    }
-
     private async Task<int> UpsertKpiValueAsync(string metricCode, DateTime bizDate, decimal? metricValue, DateTime now, S7MdpRefreshOption option)
     private async Task<int> UpsertKpiValueAsync(string metricCode, DateTime bizDate, decimal? metricValue, DateTime now, S7MdpRefreshOption option)
     {
     {
         // 沿用 S3 UpsertS3KpiValueAsync 范式:先查现存行 → UPDATE;不存在 → SELECT MAX(id)+1 显式生成 id 后 INSERT。
         // 沿用 S3 UpsertS3KpiValueAsync 范式:先查现存行 → UPDATE;不存在 → SELECT MAX(id)+1 显式生成 id 后 INSERT。

+ 58 - 10
server/Plugins/Admin.NET.Plugin.AiDOP/Job/S5MdpRefreshJob.cs

@@ -1,3 +1,6 @@
+using Admin.NET.Plugin.AiDOP.DataPlatform;
+using Admin.NET.Plugin.AiDOP.FinishedWarehouse;
+using Admin.NET.Plugin.AiDOP.Manufacturing;
 using Admin.NET.Plugin.AiDOP.MaterialWarehouse;
 using Admin.NET.Plugin.AiDOP.MaterialWarehouse;
 using Furion.Schedule;
 using Furion.Schedule;
 using Microsoft.Extensions.DependencyInjection;
 using Microsoft.Extensions.DependencyInjection;
@@ -7,12 +10,21 @@ using System.Text.Json;
 namespace Admin.NET.Plugin.AiDOP.Job;
 namespace Admin.NET.Plugin.AiDOP.Job;
 
 
 /// <summary>
 /// <summary>
-/// S5 物料仓储 T8 KPI 自动跑批任务。
-/// 每日 02:00 / 07:00 / 12:00 / 17:00 / 22:00 共 5 次(间隔 5h;22:00 是当天最后一次跑批,22:01~01:59 安静期)。
-/// 调用 S5MdpSyncTransformService.RunFullAsync(triggerType="AUTO");失败通知由 service 内部 MarkTransformRunFailedAsync 触发。
+/// T8 基础数据 + S5/S6/S7 KPI 统一刷新编排任务(方案④ owner)。
+/// 一个 tick 内:先执行 1 次 T8 基表双模式入站(源→mdp_stg_t8_*→mdp_std_t8_*),
+/// 再顺序计算 S5/S6/S7 三个模块的 T8 KPI(均读同一份本轮刷新后的 std,保证强一致)。
+/// JobId 保留 job_s5_t8_kpi_refresh(不迁移 Job 身份/历史);原 S6/S7 独立 Job 已合并至本 Job。
+/// 每日 02:00 / 07:00 / 12:00 / 17:00 / 22:00 共 5 次。
+///
+/// 失败语义(关键,勿机械照抄 Bootstrap):
+///   ① T8 inbound 是整轮硬门——失败则记错误日志并立即结束本轮,
+///      禁止继续算 S5/S6/S7(不允许退化为用上一周期 std 静默计算的 eventual consistency);
+///   ② inbound 成功后,S5→S6→S7 采用步级失败隔离(某模块失败记错误、继续下一模块);
+///   ③ OperationCanceledException 保持取消语义直接上抛,不吞。
+/// 手工补偿入口 /t8-base-mdp/inbound 仍独立保留。
 /// </summary>
 /// </summary>
 [JobDetail("job_s5_t8_kpi_refresh",
 [JobDetail("job_s5_t8_kpi_refresh",
-    Description = "S5 T8 KPI 自动跑批(5 次/天:02/07/12/17/22)",
+    Description = "T8 基础数据 + S5/S6/S7 KPI 统一刷新编排(inbound 硬门 → S5→S6→S7;5 次/天:02/07/12/17/22)",
     GroupName = "default",
     GroupName = "default",
     Concurrent = false)]
     Concurrent = false)]
 [Cron("0 2,7,12,17,22 * * *",
 [Cron("0 2,7,12,17,22 * * *",
@@ -30,22 +42,58 @@ public class S5MdpRefreshJob : IJob
     }
     }
 
 
     public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
     public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
+    {
+        // ① 硬门:T8 基表入站(源→stg→std)必须先成功;失败则本轮结束,不继续下游、不降级用旧 std。
+        try
+        {
+            using var inboundScope = _scopeFactory.CreateScope();
+            var inbound = inboundScope.ServiceProvider.GetRequiredService<T8BaseInboundMdpSyncService>();
+            var inboundResult = await inbound.RunInboundAsync(
+                tenantId: 0, fullRefresh: true, entityCode: null, cancellationToken: stoppingToken);
+            _logger.LogInformation("S5MdpRefreshJob T8 inbound 完成 {Payload}", JsonSerializer.Serialize(inboundResult));
+        }
+        catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
+        {
+            _logger.LogInformation("S5MdpRefreshJob T8 inbound 收到停止信号,本轮结束");
+            throw;
+        }
+        catch (Exception ex)
+        {
+            _logger.LogError(ex, "S5MdpRefreshJob T8 inbound 失败,本轮跳过 S5/S6/S7 KPI 刷新(不使用上一周期 std 继续计算)");
+            return;
+        }
+
+        // ② 下游:inbound 成功后,S5→S6→S7 顺序执行,步级失败隔离(各自独立 DI scope)。
+        await RunStepAsync("S5", stoppingToken, async sp =>
+            (object)await sp.GetRequiredService<S5MdpSyncTransformService>().RunFullAsync(stoppingToken, "AUTO"));
+        await RunStepAsync("S6", stoppingToken, async sp =>
+            (object)await sp.GetRequiredService<S6MdpSyncTransformService>().RunFullAsync(stoppingToken, "AUTO"));
+        await RunStepAsync("S7", stoppingToken, async sp =>
+            (object)await sp.GetRequiredService<S7MdpSyncTransformService>().RunFullAsync(stoppingToken, "AUTO"));
+    }
+
+    /// <summary>
+    /// 下游模块步级执行:独立 DI scope + 失败隔离(仿 SmartOpsKpiMdpBootstrapJob.RunStepAsync)。
+    /// 仅用于 inbound 成功之后的 S5/S6/S7;不得用于包裹 inbound(inbound 失败必须中断本轮)。
+    /// </summary>
+    private async Task RunStepAsync(string step, CancellationToken stoppingToken, Func<IServiceProvider, Task<object>> action)
     {
     {
         using var scope = _scopeFactory.CreateScope();
         using var scope = _scopeFactory.CreateScope();
-        var service = scope.ServiceProvider.GetRequiredService<S5MdpSyncTransformService>();
         try
         try
         {
         {
-            var result = await service.RunFullAsync(stoppingToken, "AUTO");
-            _logger.LogInformation("S5MdpRefreshJob 完成 {Payload}", JsonSerializer.Serialize(result));
+            stoppingToken.ThrowIfCancellationRequested();
+            var payload = await action(scope.ServiceProvider);
+            _logger.LogInformation("S5MdpRefreshJob {Step} 完成 {Payload}", step, JsonSerializer.Serialize(payload));
         }
         }
         catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
         catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
         {
         {
-            _logger.LogInformation("S5MdpRefreshJob 收到停止信号,结束本轮");
+            _logger.LogInformation("S5MdpRefreshJob {Step} 收到停止信号", step);
+            throw;
         }
         }
         catch (Exception ex)
         catch (Exception ex)
         {
         {
-            // 失败通知由 service 内部 MarkTransformRunFailedAsync 触发;此处仅记录 Job 层失败
-            _logger.LogError(ex, "S5MdpRefreshJob 执行失败");
+            // 失败通知由 service 内部 MarkTransformRunFailedAsync 触发;此处记录并继续下一步
+            _logger.LogError(ex, "S5MdpRefreshJob {Step} 执行失败,继续后续步骤", step);
         }
         }
     }
     }
 }
 }

+ 0 - 51
server/Plugins/Admin.NET.Plugin.AiDOP/Job/S6MdpRefreshJob.cs

@@ -1,51 +0,0 @@
-using Admin.NET.Plugin.AiDOP.Manufacturing;
-using Furion.Schedule;
-using Microsoft.Extensions.DependencyInjection;
-using Microsoft.Extensions.Logging;
-using System.Text.Json;
-
-namespace Admin.NET.Plugin.AiDOP.Job;
-
-/// <summary>
-/// S6 生产执行 T8 KPI 自动跑批任务。
-/// 每日 02:00 / 07:00 / 12:00 / 17:00 / 22:00 共 5 次(间隔 5h;22:00 是当天最后一次跑批,22:01~01:59 安静期)。
-/// 调用 S6MdpSyncTransformService.RunFullAsync(triggerType="AUTO");失败通知由 service 内部 MarkTransformRunFailedAsync 触发。
-/// </summary>
-[JobDetail("job_s6_t8_kpi_refresh",
-    Description = "S6 T8 KPI 自动跑批(5 次/天:02/07/12/17/22)",
-    GroupName = "default",
-    Concurrent = false)]
-[Cron("0 2,7,12,17,22 * * *",
-    TriggerId = "trigger_s6_t8_kpi_refresh",
-    Description = "每日 02:00/07:00/12:00/17:00/22:00 触发(5 字段:分 时 日 月 周,默认 CronStringFormat.Default)")]
-public class S6MdpRefreshJob : IJob
-{
-    private readonly IServiceScopeFactory _scopeFactory;
-    private readonly ILogger _logger;
-
-    public S6MdpRefreshJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
-    {
-        _scopeFactory = scopeFactory;
-        _logger = loggerFactory.CreateLogger(nameof(S6MdpRefreshJob));
-    }
-
-    public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
-    {
-        using var scope = _scopeFactory.CreateScope();
-        var service = scope.ServiceProvider.GetRequiredService<S6MdpSyncTransformService>();
-        try
-        {
-            var result = await service.RunFullAsync(stoppingToken, "AUTO");
-            _logger.LogInformation("S6MdpRefreshJob 完成 {Payload}", JsonSerializer.Serialize(result));
-        }
-        catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
-        {
-            _logger.LogInformation("S6MdpRefreshJob 收到停止信号,结束本轮");
-        }
-        catch (Exception ex)
-        {
-            // 失败通知由 service 内部 MarkTransformRunFailedAsync 触发;此处仅记录 Job 层失败
-            _logger.LogError(ex, "S6MdpRefreshJob 执行失败");
-        }
-    }
-}

+ 0 - 51
server/Plugins/Admin.NET.Plugin.AiDOP/Job/S7MdpRefreshJob.cs

@@ -1,51 +0,0 @@
-using Admin.NET.Plugin.AiDOP.FinishedWarehouse;
-using Furion.Schedule;
-using Microsoft.Extensions.DependencyInjection;
-using Microsoft.Extensions.Logging;
-using System.Text.Json;
-
-namespace Admin.NET.Plugin.AiDOP.Job;
-
-/// <summary>
-/// S7 成品仓储 T8 KPI 自动跑批任务。
-/// 每日 02:00 / 07:00 / 12:00 / 17:00 / 22:00 共 5 次(间隔 5h;22:00 是当天最后一次跑批,22:01~01:59 安静期)。
-/// 调用 S7MdpSyncTransformService.RunFullAsync(triggerType="AUTO");失败通知由 service 内部 MarkTransformRunFailedAsync 触发。
-/// </summary>
-[JobDetail("job_s7_t8_kpi_refresh",
-    Description = "S7 T8 KPI 自动跑批(5 次/天:02/07/12/17/22)",
-    GroupName = "default",
-    Concurrent = false)]
-[Cron("0 2,7,12,17,22 * * *",
-    TriggerId = "trigger_s7_t8_kpi_refresh",
-    Description = "每日 02:00/07:00/12:00/17:00/22:00 触发(5 字段:分 时 日 月 周,默认 CronStringFormat.Default)")]
-public class S7MdpRefreshJob : IJob
-{
-    private readonly IServiceScopeFactory _scopeFactory;
-    private readonly ILogger _logger;
-
-    public S7MdpRefreshJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
-    {
-        _scopeFactory = scopeFactory;
-        _logger = loggerFactory.CreateLogger(nameof(S7MdpRefreshJob));
-    }
-
-    public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
-    {
-        using var scope = _scopeFactory.CreateScope();
-        var service = scope.ServiceProvider.GetRequiredService<S7MdpSyncTransformService>();
-        try
-        {
-            var result = await service.RunFullAsync(stoppingToken, "AUTO");
-            _logger.LogInformation("S7MdpRefreshJob 完成 {Payload}", JsonSerializer.Serialize(result));
-        }
-        catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
-        {
-            _logger.LogInformation("S7MdpRefreshJob 收到停止信号,结束本轮");
-        }
-        catch (Exception ex)
-        {
-            // 失败通知由 service 内部 MarkTransformRunFailedAsync 触发;此处仅记录 Job 层失败
-            _logger.LogError(ex, "S7MdpRefreshJob 执行失败");
-        }
-    }
-}

+ 15 - 20
server/Plugins/Admin.NET.Plugin.AiDOP/Manufacturing/S6MdpSyncTransformService.cs

@@ -6,8 +6,9 @@ using System.Text.Json;
 namespace Admin.NET.Plugin.AiDOP.Manufacturing;
 namespace Admin.NET.Plugin.AiDOP.Manufacturing;
 
 
 /// <summary>
 /// <summary>
-/// S6 生产执行 — T8 KPI 数据底座与刷新转换服务。
-/// 路径 A:方老师 v5.4 KPI 字段对照表 J 列 SQL 原逻辑直发 T8 SQL Server(ConfigId=t8_v5)。
+/// S6 生产执行 — KPI 计算与刷新转换服务。
+/// 双模式:读本地标准层 mdp_std_t8_kc_*(由 T8BaseInboundMdpSyncService 从 T8 贴源→标准),不再直连 T8。
+/// 计算逻辑沿用方老师 v5.4 KPI J 列口径(等价改写为 MySQL 读 std)。
 /// 包含 KPI:S6_L1_001 工单制造满足率 / S6_L1_002 工单制造人效。
 /// 包含 KPI:S6_L1_001 工单制造满足率 / S6_L1_002 工单制造人效。
 /// </summary>
 /// </summary>
 public class S6MdpSyncTransformService : ITransient
 public class S6MdpSyncTransformService : ITransient
@@ -17,7 +18,6 @@ public class S6MdpSyncTransformService : ITransient
     private readonly ILogger<S6MdpSyncTransformService> _logger;
     private readonly ILogger<S6MdpSyncTransformService> _logger;
     private const string JobCode = "S6_MDP_SYNC_TRANSFORM";
     private const string JobCode = "S6_MDP_SYNC_TRANSFORM";
     private const string JobName = "S6 生产执行 MDP 同步与转换";
     private const string JobName = "S6 生产执行 MDP 同步与转换";
-    private const string T8ConfigId = "t8_v5";
     private const string ModuleCode = "S6";
     private const string ModuleCode = "S6";
     // FAILURE-NOTIFICATION-1:超级管理员 superAdmin.NET(AccountType=999)
     // FAILURE-NOTIFICATION-1:超级管理员 superAdmin.NET(AccountType=999)
     private const long NoticeReceiverUserId = 1300000000101L;
     private const long NoticeReceiverUserId = 1300000000101L;
@@ -90,21 +90,22 @@ public class S6MdpSyncTransformService : ITransient
     {
     {
         var sub = new KpiBuildSubResult();
         var sub = new KpiBuildSubResult();
 
 
+        // 双模式:读本地标准层 mdp_std_t8_*(源 identity Id→src_id),语义等价于原直发 T8 SQL。
         const string sql = @"
         const string sql = @"
 select a.noid as noid, b.rwnoid as rwnoid, b.code as code, b.sl as sl, sum(d.slzx) as slzx
 select a.noid as noid, b.rwnoid as rwnoid, b.code as code, b.sl as sl, sum(d.slzx) as slzx
- from kc_dd_head a with(nolock)
- left join kc_dd_list b with(nolock) on a.Id=b.idid
+ from mdp_std_t8_kc_dd_head a
+ left join mdp_std_t8_kc_dd_list b on a.src_id=b.idid
  left join (
  left join (
    select b.lynoid as lynoid, b.code as code,
    select b.lynoid as lynoid, b.code as code,
-          convert(varchar(10), a.shtime, 23) as shtime, b.slzx as slzx
-    from kc_tz_head a with(nolock)
-    inner join kc_tz_list b with(nolock) on a.Id=b.idid
+          date(a.shtime) as shtime, b.slzx as slzx
+    from mdp_std_t8_kc_tz_head a
+    inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
     where a.ztid=@ztid and a.lbs='生产入库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
     where a.ztid=@ztid and a.lbs='生产入库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  ) d on b.rwnoid=d.lynoid and b.code=d.code
  ) d on b.rwnoid=d.lynoid and b.code=d.code
  where a.ztid=@ztid and a.lbs='生产任务' and a.zf=0 and a.shyn=1 and d.shtime<=b.jhdate
  where a.ztid=@ztid and a.lbs='生产任务' and a.zf=0 and a.shyn=1 and d.shtime<=b.jhdate
  group by a.noid, b.rwnoid, b.code, b.sl";
  group by a.noid, b.rwnoid, b.code, b.sl";
 
 
-        var rows = await QueryT8Async<S6MfgFulfillmentRow>(sql, new[] { new SugarParameter("@ztid", option.SourceZtid) });
+        var rows = await _db.Ado.SqlQueryAsync<S6MfgFulfillmentRow>(sql, new[] { new SugarParameter("@ztid", option.SourceZtid) });
         sub.T8Rows = rows.Count;
         sub.T8Rows = rows.Count;
 
 
         var dwdAffected = 0;
         var dwdAffected = 0;
@@ -160,14 +161,14 @@ ON DUPLICATE KEY UPDATE
 
 
         const string sqlNumer = @"
         const string sqlNumer = @"
 select count(*) as ddnum
 select count(*) as ddnum
- from kc_tz_head a with(nolock)
- inner join kc_tz_list b with(nolock) on a.Id=b.idid
+ from mdp_std_t8_kc_tz_head a
+ inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
  where a.ztid=@ztid and a.lbs='生产入库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  where a.ztid=@ztid and a.lbs='生产入库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
   and a.date0 between @startDate and @endDate
   and a.date0 between @startDate and @endDate
   and b.slzx>0 and (b.slzx>=b.sl or b.gdyn=1)";
   and b.slzx>0 and (b.slzx>=b.sl or b.gdyn=1)";
         const string sqlDenom = @"
         const string sqlDenom = @"
 select count(*) as penum
 select count(*) as penum
- from sys_pelist with(nolock)
+ from mdp_std_t8_sys_pelist
  where ztid=@ztid and zzzt='在职' and gw='生产'";
  where ztid=@ztid and zzzt='在职' and gw='生产'";
 
 
         var pNumer = new[]
         var pNumer = new[]
@@ -178,8 +179,8 @@ select count(*) as penum
         };
         };
         var pDenom = new[] { new SugarParameter("@ztid", option.SourceZtid) };
         var pDenom = new[] { new SugarParameter("@ztid", option.SourceZtid) };
 
 
-        var numerRows = await QueryT8Async<S6CountRow>(sqlNumer, pNumer);
-        var denomRows = await QueryT8Async<S6PeNumRow>(sqlDenom, pDenom);
+        var numerRows = await _db.Ado.SqlQueryAsync<S6CountRow>(sqlNumer, pNumer);
+        var denomRows = await _db.Ado.SqlQueryAsync<S6PeNumRow>(sqlDenom, pDenom);
         sub.T8Rows = numerRows.Count + denomRows.Count;
         sub.T8Rows = numerRows.Count + denomRows.Count;
 
 
         int? doneCount = numerRows.FirstOrDefault()?.ddnum;
         int? doneCount = numerRows.FirstOrDefault()?.ddnum;
@@ -230,12 +231,6 @@ ON DUPLICATE KEY UPDATE
 
 
     // ─────────────────────────────────────────────────────────────────────────
     // ─────────────────────────────────────────────────────────────────────────
 
 
-    private async Task<List<T>> QueryT8Async<T>(string sql, SugarParameter[] parameters)
-    {
-        var t8 = _db.AsTenant().GetConnectionScope(T8ConfigId);
-        return await t8.Ado.SqlQueryAsync<T>(sql, parameters);
-    }
-
     private async Task<int> UpsertKpiValueAsync(string metricCode, DateTime bizDate, decimal? metricValue, DateTime now, S6MdpRefreshOption option)
     private async Task<int> UpsertKpiValueAsync(string metricCode, DateTime bizDate, decimal? metricValue, DateTime now, S6MdpRefreshOption option)
     {
     {
         // 沿用 S3 UpsertS3KpiValueAsync 范式:先查现存行 → UPDATE;不存在 → SELECT MAX(id)+1 显式生成 id 后 INSERT。
         // 沿用 S3 UpsertS3KpiValueAsync 范式:先查现存行 → UPDATE;不存在 → SELECT MAX(id)+1 显式生成 id 后 INSERT。

+ 20 - 16
server/Plugins/Admin.NET.Plugin.AiDOP/MaterialWarehouse/S5MdpSyncTransformService.cs

@@ -6,11 +6,12 @@ using System.Text.Json;
 namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
 namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
 
 
 /// <summary>
 /// <summary>
-/// S5 物料仓储 — T8 KPI 数据底座与刷新转换服务。
-/// 路径 A:方老师 v5.4 KPI 字段对照表 J 列 SQL 原逻辑直发 T8 SQL Server(ConfigId=t8_v5);
-/// 一期不做 mdp_stg_t8_* 贴源层;结果直接落 dwd_t8_* 与 ado_s9_kpi_value_l1_day。
-/// 包含 KPI:S5_L1_001 物料上线周期 / S5_L1_002 物料上线满足率 /
-///          S5_L1_003 物料仓储人效 / S5_L1_004 品类物料库存周转。
+/// S5 物料仓储 — KPI 计算与刷新转换服务。边界:7 迁 2 留。
+/// 已中台化(读本地标准层 mdp_std_t8_*,由 T8BaseInboundMdpSyncService 从 T8 贴源→标准):
+///   S5_L1_001 物料上线周期 / S5_L1_003 物料仓储人效。
+/// 仍 legacy 直连 T8(QueryT8Async,ConfigId=t8_v5):
+///   S5_L1_002 物料上线满足率(依赖 T8 报表 Cj_Bg_Head_Rep)/ S5_L1_004 品类物料库存周转(依赖 T8 TVF Rep_总账_存货_V3)。
+/// 结果统一落 dwd_t8_* 与 ado_s9_kpi_value_l1_day。计算口径沿用方老师 v5.4 KPI J 列。
 /// </summary>
 /// </summary>
 public class S5MdpSyncTransformService : ITransient
 public class S5MdpSyncTransformService : ITransient
 {
 {
@@ -103,22 +104,23 @@ public class S5MdpSyncTransformService : ITransient
     {
     {
         var sub = new KpiBuildSubResult();
         var sub = new KpiBuildSubResult();
 
 
+        // 双模式:读本地标准层 mdp_std_t8_*(源 identity Id→src_id),语义等价于原直发 T8 SQL。
         const string sqlOnline = @"
         const string sqlOnline = @"
 select b.code as code, min(a.shtime) as shtime
 select b.code as code, min(a.shtime) as shtime
- from kc_tz_head a with(nolock)
- inner join kc_tz_list b with(nolock) on a.Id=b.idid
+ from mdp_std_t8_kc_tz_head a
+ inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
  where a.ztid=@ztid and a.lbs='生产领料' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  where a.ztid=@ztid and a.lbs='生产领料' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  group by b.code";
  group by b.code";
         const string sqlReceipt = @"
         const string sqlReceipt = @"
 select b.code as code, min(a.shtime) as shtime
 select b.code as code, min(a.shtime) as shtime
- from kc_tz_head a with(nolock)
- inner join kc_tz_list b with(nolock) on a.Id=b.idid
+ from mdp_std_t8_kc_tz_head a
+ inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
  where a.ztid=@ztid and a.lbs='采购入库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  where a.ztid=@ztid and a.lbs='采购入库' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  group by b.code";
  group by b.code";
 
 
         var p = new[] { new SugarParameter("@ztid", option.SourceZtid) };
         var p = new[] { new SugarParameter("@ztid", option.SourceZtid) };
-        var onlineRows = await QueryT8Async<S5OnlineCycleRow>(sqlOnline, p);
-        var receiptRows = await QueryT8Async<S5OnlineCycleRow>(sqlReceipt, p);
+        var onlineRows = await _db.Ado.SqlQueryAsync<S5OnlineCycleRow>(sqlOnline, p);
+        var receiptRows = await _db.Ado.SqlQueryAsync<S5OnlineCycleRow>(sqlReceipt, p);
         sub.T8Rows = onlineRows.Count + receiptRows.Count;
         sub.T8Rows = onlineRows.Count + receiptRows.Count;
 
 
         var onlineByCode = onlineRows.Where(r => !string.IsNullOrEmpty(r.code))
         var onlineByCode = onlineRows.Where(r => !string.IsNullOrEmpty(r.code))
@@ -261,13 +263,13 @@ ON DUPLICATE KEY UPDATE
 
 
         const string sqlNumer = @"
         const string sqlNumer = @"
 select sum(b.slzx) as slzx
 select sum(b.slzx) as slzx
- from kc_tz_head a with(nolock)
- inner join kc_tz_list b with(nolock) on a.Id=b.idid
+ from mdp_std_t8_kc_tz_head a
+ inner join mdp_std_t8_kc_tz_list b on a.src_id=b.idid
  where a.ztid=@ztid and a.lbs='生产领料' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
  where a.ztid=@ztid and a.lbs='生产领料' and a.hzyn=0 and a.zfyn=0 and a.shyn=1
   and a.shtime between @startDate and @endDate";
   and a.shtime between @startDate and @endDate";
         const string sqlDenom = @"
         const string sqlDenom = @"
 select count(*) as penum
 select count(*) as penum
- from sys_pelist with(nolock)
+ from mdp_std_t8_sys_pelist
  where ztid=@ztid and zzzt='在职' and gw='仓管'";
  where ztid=@ztid and zzzt='在职' and gw='仓管'";
 
 
         var pNumer = new[]
         var pNumer = new[]
@@ -278,8 +280,8 @@ select count(*) as penum
         };
         };
         var pDenom = new[] { new SugarParameter("@ztid", option.SourceZtid) };
         var pDenom = new[] { new SugarParameter("@ztid", option.SourceZtid) };
 
 
-        var numerRows = await QueryT8Async<S5SumQtyRow>(sqlNumer, pNumer);
-        var denomRows = await QueryT8Async<S5CountRow>(sqlDenom, pDenom);
+        var numerRows = await _db.Ado.SqlQueryAsync<S5SumQtyRow>(sqlNumer, pNumer);
+        var denomRows = await _db.Ado.SqlQueryAsync<S5CountRow>(sqlDenom, pDenom);
         sub.T8Rows = numerRows.Count + denomRows.Count;
         sub.T8Rows = numerRows.Count + denomRows.Count;
 
 
         decimal? onlineQty = numerRows.FirstOrDefault()?.slzx;
         decimal? onlineQty = numerRows.FirstOrDefault()?.slzx;
@@ -420,6 +422,8 @@ ON DUPLICATE KEY UPDATE
     // 跨库 / 写入 / 日志 封装
     // 跨库 / 写入 / 日志 封装
     // ─────────────────────────────────────────────────────────────────────────
     // ─────────────────────────────────────────────────────────────────────────
 
 
+    // legacy 直连 T8:仅 S5_L1_002(Cj_Bg_Head_Rep 报表)/ S5_L1_004(Rep_总账_存货_V3 TVF)仍用;
+    // 这两类源非干净基表,暂不纳入双模式贴源,待底表摸清后单独立批。已中台化的 001/003 不再走此方法。
     private async Task<List<T>> QueryT8Async<T>(string sql, SugarParameter[] parameters)
     private async Task<List<T>> QueryT8Async<T>(string sql, SugarParameter[] parameters)
     {
     {
         var t8 = _db.AsTenant().GetConnectionScope(T8ConfigId);
         var t8 = _db.AsTenant().GetConnectionScope(T8ConfigId);