Pārlūkot izejas kodu

feat(mdp): S5生产退料 bespoke→B-①通用管线 + dopdemorq 双源只读入站(server 1.0.271)

- ProductionReturnMdpSyncService: 直读 NbrMaster/NbrDetail 改为执行器灌 mdp_stg_production_return(头/明细共表)→ transform 读 stg(MdpJsonSql 跨源类型兼容)→ mdp_std_production_return(_detail)
- 退料专属: IsReturn 判假、RecID↔NbrRecID 关联、tenant_id 从 raw_data 保留(本库租户不回归)、双 std 专表各一次 FULL Replace(头先明细后回填 std_head_id)
- 1.0.271.sql: 登记 4 实体(本库 MASTER/DETAIL status=1,DOPDEMORQ_SQLSERVER MASTER/DETAIL status=0 就位不启用),全 FULL
- 前端/核心/DTO 零改;本库 WOD 冒烟 1头1明细等价、源表不写、类型与 std_head_id 回填全过
YY968XX 1 nedēļu atpakaļ
vecāks
revīzija
72213702a1

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

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

+ 52 - 0
server/Admin.NET.Web.Entry/UpdateScripts/1.0.271.sql

@@ -0,0 +1,52 @@
+-- 1.0.271:S5 生产退料 B-① 迁通用管线 + dopdemorq 双源入站实体登记(本库 status=1 / SQLSERVER status=0)
+--
+-- 背景:ProductionReturn 由 bespoke 直读 NbrMaster/NbrDetail 迁入通用管线
+--       (源→mdp_stg_production_return[头/明细共表,source_table 区分]→transform→mdp_std_production_return(_detail))。
+--       登记 4 个入站实体(本库 master+detail、dopdemorq SQL Server master+detail),均落 mdp_stg_production_return、FULL。
+-- 状态:正式脚本(1.0.271,已注册 csproj Copy,随 AutoVersionUpdate 执行)。
+--
+-- 硬约束:
+--   1) 本库实体 S5_PRODUCTION_RETURN_MASTER/_DETAIL status=1:RunFullAsync 迁 stg 管线后依赖它(否则 s5-production-return-mdp/refresh 端点断)。
+--   2) SQLSERVER 实体 status=0:就位不启用,合入后不改变现有运行行为;启用需带外注密码+受控切源。
+--   3) 四实体均 sync_mode='FULL'(本批 FULL-only,规避增量边界;不启自动增量)。
+--   4) source_table_name='NbrMaster'/'NbrDetail'(两库同名);transform 按 Type='WOD' AND IsActive(判真) AND IsReturn(判假) 过滤(退料)。
+--      tenant_id 由入站上下文赋值,本库源含该列(raw 保留),SQL Server 源无此列(回落上下文)。
+-- 幂等:WHERE NOT EXISTS 守卫,可重复执行。
+
+-- 本库入站实体(AIDOPDEV_MYSQL, status=1)
+INSERT INTO mdp_entity
+  (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name,
+   sync_mode, batch_size, biz_key_expr, status, remark, create_time, update_time)
+SELECT 0, (SELECT id FROM mdp_source WHERE source_code='AIDOPDEV_MYSQL'),
+       'S5_PRODUCTION_RETURN_MASTER', 'S5生产退料主入站(本库)', 'TABLE',
+       'NbrMaster', 'mdp_stg_production_return',
+       'FULL', 5000, 'Domain,Nbr', 1, 'S5生产退料 B-① 本库头入站', NOW(), NOW()
+WHERE NOT EXISTS (SELECT 1 FROM mdp_entity WHERE entity_code='S5_PRODUCTION_RETURN_MASTER');
+
+INSERT INTO mdp_entity
+  (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name,
+   sync_mode, batch_size, biz_key_expr, status, remark, create_time, update_time)
+SELECT 0, (SELECT id FROM mdp_source WHERE source_code='AIDOPDEV_MYSQL'),
+       'S5_PRODUCTION_RETURN_DETAIL', 'S5生产退料明细入站(本库)', 'TABLE',
+       'NbrDetail', 'mdp_stg_production_return',
+       'FULL', 5000, 'Domain,Nbr,Line', 1, 'S5生产退料 B-① 本库明细入站', NOW(), NOW()
+WHERE NOT EXISTS (SELECT 1 FROM mdp_entity WHERE entity_code='S5_PRODUCTION_RETURN_DETAIL');
+
+-- dopdemorq SQL Server 入站实体(DOPDEMORQ_SQLSERVER, status=0 就位不启用)
+INSERT INTO mdp_entity
+  (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name,
+   sync_mode, batch_size, biz_key_expr, status, remark, create_time, update_time)
+SELECT 0, (SELECT id FROM mdp_source WHERE source_code='DOPDEMORQ_SQLSERVER'),
+       'S5_PRODUCTION_RETURN_MASTER_SQLSERVER', 'S5生产退料主入站(SQLServer)', 'TABLE',
+       'NbrMaster', 'mdp_stg_production_return',
+       'FULL', 5000, 'Domain,Nbr', 0, 'S5生产退料 B-① 双源 disabled', NOW(), NOW()
+WHERE NOT EXISTS (SELECT 1 FROM mdp_entity WHERE entity_code='S5_PRODUCTION_RETURN_MASTER_SQLSERVER');
+
+INSERT INTO mdp_entity
+  (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name,
+   sync_mode, batch_size, biz_key_expr, status, remark, create_time, update_time)
+SELECT 0, (SELECT id FROM mdp_source WHERE source_code='DOPDEMORQ_SQLSERVER'),
+       'S5_PRODUCTION_RETURN_DETAIL_SQLSERVER', 'S5生产退料明细入站(SQLServer)', 'TABLE',
+       'NbrDetail', 'mdp_stg_production_return',
+       'FULL', 5000, 'Domain,Nbr,Line', 0, 'S5生产退料 B-① 双源 disabled', NOW(), NOW()
+WHERE NOT EXISTS (SELECT 1 FROM mdp_entity WHERE entity_code='S5_PRODUCTION_RETURN_DETAIL_SQLSERVER');

+ 286 - 61
server/Plugins/Admin.NET.Plugin.AiDOP/MaterialWarehouse/ProductionReturnMdpSyncService.cs

@@ -1,37 +1,53 @@
+using Admin.NET.Plugin.AiDOP.DataPlatform;
+using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
+
 namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
 
 /// <summary>
-/// S5 生产退料单 数据中台只读同步转换服务(DOP 内部,head-detail 双表,独立于 S5MdpSyncTransformService 的 KPI 管线)。
+/// S5 生产退料单 数据中台只读同步转换服务(DOP 内部,head-detail 双 std 专表,独立于 S5MdpSyncTransformService 的 KPI 管线)。
+///
+/// B-① 迁通用管线:源(NbrMaster/NbrDetail, Type='WOD') 经执行器灌 mdp_stg_production_return(头/明细共表,source_table 区分)
+///   → transform 读 stg(MdpJsonSql 跨源类型兼容)→ mdp_std_production_return(_detail)(typed,明细回填 std_head_id)。
+///   维表 DepartmentMaster(d) / ItemMaster(i) / LocationMaster(lf,lt) 仍本库 LEFT JOIN(富化 desc,业务 code 由事实表保留)。
 ///
-/// 源:aidopdev NbrMaster(m) + NbrDetail(dt),业务类型 Type='WOD';头明细关联 dt.NbrRecID -> m.RecID。
-///     维表 DepartmentMaster(d) / ItemMaster(i) / LocationMaster(lf,lt) 全 LEFT JOIN(LocationMaster 列名小写 location/descr)。
-/// 链路:NbrMaster/NbrDetail(WOD) 只读 → mdp_std_production_return(_detail)(typed)。
-///     不经 stg、不建 dwd(只读列表/详情);WOD=退料单业务类型有 S2 WorkOrderSchedulingService 代码实证。
+/// 双源:本库实体 S5_PRODUCTION_RETURN_MASTER/_DETAIL(AIDOPDEV_MYSQL, status=1);
+///   dopdemorq 第二源 S5_PRODUCTION_RETURN_MASTER_SQLSERVER/_DETAIL_SQLSERVER(DOPDEMORQ_SQLSERVER, status=0 就位不启用)。
 ///
 /// 约束:
-///   - 只读源表,仅写 mdp_std_production_return(_detail);绝不 INSERT/UPDATE/DELETE NbrMaster/NbrDetail;不碰 S2 写路径。
-///   - 头过滤 m.Type='WOD' AND m.IsActive=1 AND m.IsReturn=0;明细 dt.Type='WOD' 且经 NbrRecID join WOD 头。
-///   - status_desc 规则本批仅:Status='C' -> '关闭',其他 -> 原值(UPPER)。
-///   - 退料人 return_user 本批取 m.User1(列表 SQL 口径);user1/user2 原样落库,页面策略后置微调。
-///   - trans_type_text 恒空(旧系统口径);不做 danjia/jiage/库存事务/状态流转。
-///   - WOD 为 0 行时转换成功完成、处理数为 0,不报错。
+///   - 只读源/贴源,仅写 mdp_stg_production_return / mdp_std_production_return(_detail);绝不 INSERT/UPDATE/DELETE NbrMaster/NbrDetail;不碰 S2 写路径。
+///   - 头过滤 Type='WOD' AND IsActive(判真) AND IsReturn(判假);明细 Type='WOD' 且经 RecID↔NbrRecID join WOD 头。
+///   - status_desc 规则:Status='C' -> '关闭',其他 -> 原值(UPPER)。trans_type_text 恒空(旧系统口径);不做 danjia/库存事务/状态流转。
+///   - tenant_id 从 raw_data 保留(本库源含该列 → std 租户不回归;SQL Server 源无此列 → 回落入站上下文 tenant)。
+///   - WOD 为 0 行时转换成功完成、处理数为 0,不报错。切源 FULL Replace 头/明细各一次单事务、仅按 tenant、Pull 成功后才 destructive。
 /// </summary>
 public class ProductionReturnMdpSyncService : ITransient
 {
     private const string JobCode = "S5_PRODUCTION_RETURN_MDP_SYNC";
+    private const string InboundMasterEntityCode = "S5_PRODUCTION_RETURN_MASTER";
+    private const string InboundDetailEntityCode = "S5_PRODUCTION_RETURN_DETAIL";
+
+    private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
+    private const string SqlServerMasterEntityCode = "S5_PRODUCTION_RETURN_MASTER_SQLSERVER";
+    private const string SqlServerDetailEntityCode = "S5_PRODUCTION_RETURN_DETAIL_SQLSERVER";
+
+    private const string StdHeadTable = "mdp_std_production_return";
+    private const string StdDetailTable = "mdp_std_production_return_detail";
 
     private readonly ISqlSugarClient _db;
+    private readonly MdpSourcePullDispatcher _pullDispatcher;
 
-    public ProductionReturnMdpSyncService(ISqlSugarClient db)
+    public ProductionReturnMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher)
     {
         _db = db;
+        _pullDispatcher = pullDispatcher;
     }
 
-    /// <summary>全量同步转换:源(WOD 头/明细) → 标准层(头/明细)。</summary>
+    /// <summary>全量:本地 DB 执行器灌 stg(头+明细) → 标准层头/明细(读 stg)。</summary>
     public async Task<ProductionReturnMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
     {
         cancellationToken.ThrowIfCancellationRequested();
         await EnsureTablesAsync();
+        await EnsureStgTableAsync();
 
         var now = DateTime.Now;
         var batchId = $"S5_PROD_RTN_FULL_{now:yyyyMMddHHmmss}";
@@ -40,6 +56,14 @@ public class ProductionReturnMdpSyncService : ITransient
 
         try
         {
+            var pullCtx = new MdpPullContext
+            {
+                TenantId = 0,
+                FullRefresh = true,
+                TaskCode = "S5_PRODUCTION_RETURN_INBOUND",
+                BatchId = $"{batchId}_PULL"
+            };
+            await PopulateStgAsync(pullCtx, cancellationToken);
             result.HeadRows = await TransformHeadStandardAsync(batchId, now);
             result.DetailRows = await TransformDetailStandardAsync(batchId, now);
             await MarkRunSuccessAsync(runLogId, now, result);
@@ -52,6 +76,153 @@ public class ProductionReturnMdpSyncService : ITransient
         }
     }
 
+    /// <summary>双模式入站:执行器抽主/明细 → stg,再跑 WOD 标准层头/明细(读 stg)。</summary>
+    public async Task<ProductionReturnInboundResult> RunInboundAsync(
+        long tenantId = 0,
+        bool fullRefresh = false,
+        CancellationToken cancellationToken = default)
+    {
+        cancellationToken.ThrowIfCancellationRequested();
+        await EnsureTablesAsync();
+        await EnsureStgTableAsync();
+
+        var now = DateTime.Now;
+        var pullCtx = new MdpPullContext
+        {
+            TenantId = tenantId,
+            FullRefresh = fullRefresh,
+            TaskCode = "S5_PRODUCTION_RETURN_INBOUND",
+            BatchId = $"S5_PROD_RTN_IN_{now:yyyyMMddHHmmss}"
+        };
+        var pull = await PopulateStgAsync(pullCtx, cancellationToken);
+
+        var batchId = $"S5_PROD_RTN_STD_{now:yyyyMMddHHmmss}";
+        var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
+        var result = new ProductionReturnMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
+        try
+        {
+            result.HeadRows = await TransformHeadStandardAsync(batchId, now);
+            result.DetailRows = await TransformDetailStandardAsync(batchId, now);
+            await MarkRunSuccessAsync(runLogId, now, result);
+        }
+        catch (Exception ex)
+        {
+            await MarkRunFailedAsync(runLogId, now, ex.Message);
+            throw;
+        }
+
+        return new ProductionReturnInboundResult
+        {
+            PullBatchId = pullCtx.BatchId,
+            RowsPulled = pull.RowsPulled,
+            RowsWrittenStg = pull.RowsWritten,
+            NewCursor = pull.NewCursor,
+            TransformBatchId = batchId,
+            StdHeadRows = result.HeadRows,
+            StdDetailRows = result.DetailRows
+        };
+    }
+
+    /// <summary>
+    /// 双源切换 FULL Replace:从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 stg,
+    /// 成功后头/明细各一次单事务 FULL 重建(按 tenant 精确隔离),消除旧源独有业务键残留。
+    /// 顺序:先 head(写 std 头,供明细回填 std_head_id),再 detail。两次独立事务,Pull 成功后才 destructive。
+    /// SQLSERVER 实体默认 status=0(就位不启用);dopdemorq 六表当前空,Phase 1 不实际执行 destructive 切换。
+    /// </summary>
+    public async Task<ProductionReturnInboundResult> RunSourceSwitchFullAsync(
+        string sourceCode = SqlServerSourceCode,
+        string masterEntityCode = SqlServerMasterEntityCode,
+        string detailEntityCode = SqlServerDetailEntityCode,
+        long tenantId = 0,
+        CancellationToken cancellationToken = default)
+    {
+        cancellationToken.ThrowIfCancellationRequested();
+        await EnsureTablesAsync();
+        await EnsureStgTableAsync();
+
+        var now = DateTime.Now;
+        // 1) PullAll 抽尽 + fullRefresh=true。任一 Pull 失败会抛异常,std 未动。
+        var pullCtx = new MdpPullContext
+        {
+            TenantId = tenantId,
+            FullRefresh = true,
+            TaskCode = "S5_PRODUCTION_RETURN_INBOUND",
+            BatchId = $"S5_PROD_RTN_SW_{now:yyyyMMddHHmmss}"
+        };
+        var master = await _pullDispatcher.PullAllByEntityCodeAsync(masterEntityCode, pullCtx, cancellationToken);
+        var detail = await _pullDispatcher.PullAllByEntityCodeAsync(detailEntityCode, pullCtx, cancellationToken);
+
+        // 2) FULL Replace:头先删建(供明细回填 std_head_id),明细后删建;各自单事务、仅按 tenant、仅当前 source_system。
+        var batchId = $"S5_PROD_RTN_SWSTD_{now:yyyyMMddHHmmss}";
+        var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
+        var result = new ProductionReturnMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
+        try
+        {
+            result.HeadRows = await MdpStdFullReplace.ReplaceAsync(
+                _db, StdHeadTable, tenantId, extraWhere: null,
+                insertScopedAsync: () => TransformHeadStandardAsync(batchId, now, sourceCode),
+                cancellationToken);
+            result.DetailRows = await MdpStdFullReplace.ReplaceAsync(
+                _db, StdDetailTable, tenantId, extraWhere: null,
+                insertScopedAsync: () => TransformDetailStandardAsync(batchId, now, sourceCode),
+                cancellationToken);
+            await MarkRunSuccessAsync(runLogId, now, result);
+        }
+        catch (Exception ex)
+        {
+            await MarkRunFailedAsync(runLogId, now, ex.Message);
+            throw;
+        }
+
+        return new ProductionReturnInboundResult
+        {
+            PullBatchId = pullCtx.BatchId,
+            RowsPulled = master.RowsPulled + detail.RowsPulled,
+            RowsWrittenStg = master.RowsWritten + detail.RowsWritten,
+            NewCursor = detail.NewCursor ?? master.NewCursor,
+            TransformBatchId = batchId,
+            StdHeadRows = result.HeadRows,
+            StdDetailRows = result.DetailRows
+        };
+    }
+
+    private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
+        MdpPullContext pullCtx, CancellationToken cancellationToken)
+    {
+        var master = await _pullDispatcher.PullByEntityCodeAsync(InboundMasterEntityCode, pullCtx, cancellationToken);
+        var detail = await _pullDispatcher.PullByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
+        return (
+            master.RowsPulled + detail.RowsPulled,
+            master.RowsWritten + detail.RowsWritten,
+            detail.NewCursor ?? master.NewCursor);
+    }
+
+    private async Task EnsureStgTableAsync()
+    {
+        await _db.Ado.ExecuteCommandAsync(
+            """
+            CREATE TABLE IF NOT EXISTS mdp_stg_production_return (
+                id BIGINT PRIMARY KEY AUTO_INCREMENT,
+                tenant_id BIGINT NOT NULL,
+                source_system VARCHAR(50) NULL,
+                source_table VARCHAR(200),
+                source_row_id VARCHAR(200),
+                source_biz_key VARCHAR(300) NULL,
+                raw_data JSON,
+                sync_batch_id VARCHAR(100),
+                sync_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP,
+                process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
+                process_message VARCHAR(500) NULL,
+                create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+                update_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+                UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
+                KEY idx_batch (sync_batch_id),
+                KEY idx_src (source_table, source_row_id),
+                KEY idx_tenant (tenant_id)
+            ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='S5生产退料执行器贴源层(头/明细共表,source_table 区分)'
+            """);
+    }
+
     /// <summary>防御式建表(与 UpdateScripts/1.0.213.sql 同构,幂等)。</summary>
     private async Task EnsureTablesAsync()
     {
@@ -138,33 +309,54 @@ public class ProductionReturnMdpSyncService : ITransient
     }
 
     /// <summary>
-    /// 标准化头:NbrMaster(WOD) + DepartmentMaster -> mdp_std_production_return。返回处理头行数。
-    /// SQL 由旧系统头 SQL 翻译(SQL Server -&gt; MySQL):IsNull-&gt;IFNULL、upper() 保留、
-    /// rtrim(a+' '+b)-&gt;TRIM(CONCAT(...))、with(nolock) 去除。
+    /// 标准化头:stg(NbrMaster, WOD) + DepartmentMaster(本库) → mdp_std_production_return。返回处理头行数。
+    /// 全字段经 MdpJsonSql 跨源类型兼容;tenant_id 从 raw_data 保留(本库源含列 → 不回归;SQL Server 源无列 → 回落 stg 上下文 tenant)。
     /// </summary>
-    private async Task<int> TransformHeadStandardAsync(string batchId, DateTime now)
+    private async Task<int> TransformHeadStandardAsync(string batchId, DateTime now, string? sourceSystem = null)
     {
+        // WOD 头判定:Type='WOD' AND IsActive(判真) AND IsReturn(判假, 兼容 bit true/false 与 0/1)。
+        var srcClause = sourceSystem == null ? "" : " AND m.source_system=@Src";
+        var mTenant = $"IFNULL({MdpJsonSql.Int("m", "tenant_id")}, IFNULL(m.tenant_id, 0))";
+        var statusExpr = $"UPPER(IFNULL({MdpJsonSql.Str("m", "Status")}, ''))";
+        var headWhere =
+            $"""
+            m.source_table='NbrMaster'
+              AND {MdpJsonSql.Str("m", "Type")}='WOD'
+              AND {MdpJsonSql.BoolTrue("m", "IsActive")}
+              AND LOWER({MdpJsonSql.Raw("m", "IsReturn")}) IN ('0','false'){srcClause}
+            """;
+
+        var countPars = new List<SugarParameter>();
+        if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
         var rows = await _db.Ado.GetIntAsync(
-            "SELECT COUNT(1) FROM NbrMaster m WHERE m.Type='WOD' AND m.IsActive=1 AND m.IsReturn=0");
-        await _db.Ado.ExecuteCommandAsync(
-            """
+            $"SELECT COUNT(1) FROM mdp_stg_production_return m WHERE {headWhere}", countPars);
+
+        var insertSql =
+            $"""
             INSERT INTO mdp_std_production_return
             (tenant_id, factory_id, source_system, domain, rec_id, nbr, return_date, status, status_desc,
              work_ord, department, department_desc, qty_ord, prod_line, applicant_name, return_user,
              user1, user2, remark, create_user, source_create_time, trans_type, trans_type_text,
              eff_date, ufld1, source_biz_key, sync_batch_id, sync_time)
             SELECT
-                IFNULL(m.tenant_id, 0), 1, 'AIDOP', m.Domain, m.RecID, m.Nbr, m.Date,
-                UPPER(IFNULL(m.Status, '')),
-                CASE WHEN UPPER(IFNULL(m.Status,''))='C' THEN '关闭' ELSE UPPER(IFNULL(m.Status,'')) END,
-                m.WorkOrd, m.Department, TRIM(CONCAT(m.Department, ' ', IFNULL(d.Descr, ''))),
-                m.QtyOrd, m.ProdLine, m.Name, CAST(m.User1 AS CHAR),
-                CAST(m.User1 AS CHAR), CAST(m.User2 AS CHAR), m.Remark, m.CreateUser, m.CreateTime,
-                m.TransType, '', m.EffDate, m.Ufld1,
-                CONCAT(m.Domain, ':', m.RecID), @BatchId, @Now
-            FROM NbrMaster m
-            LEFT JOIN DepartmentMaster d ON d.Domain = m.Domain AND d.Department = m.Department
-            WHERE m.Type='WOD' AND m.IsActive=1 AND m.IsReturn=0
+                {mTenant}, 1, IFNULL(NULLIF(m.source_system,''), 'AIDOP'),
+                {MdpJsonSql.Str("m", "Domain")}, {MdpJsonSql.Int("m", "RecID")}, {MdpJsonSql.Str("m", "Nbr")},
+                {MdpJsonSql.DateTimeSec("m", "Date")},
+                {statusExpr},
+                CASE WHEN {statusExpr}='C' THEN '关闭' ELSE {statusExpr} END,
+                {MdpJsonSql.Str("m", "WorkOrd")}, {MdpJsonSql.Str("m", "Department")},
+                TRIM(CONCAT(IFNULL({MdpJsonSql.Str("m", "Department")},''), ' ', IFNULL(d.Descr, ''))),
+                {MdpJsonSql.Dec("m", "QtyOrd", 18, 5)}, {MdpJsonSql.Str("m", "ProdLine")}, {MdpJsonSql.Str("m", "Name")},
+                {MdpJsonSql.Str("m", "User1")},
+                {MdpJsonSql.Str("m", "User1")}, {MdpJsonSql.Str("m", "User2")}, {MdpJsonSql.Str("m", "Remark")},
+                {MdpJsonSql.Str("m", "CreateUser")}, {MdpJsonSql.DateTimeSec("m", "CreateTime")},
+                {MdpJsonSql.Str("m", "TransType")}, '', {MdpJsonSql.DateTimeSec("m", "EffDate")}, {MdpJsonSql.Str("m", "Ufld1")},
+                CONCAT(IFNULL({MdpJsonSql.Str("m", "Domain")},''), ':', IFNULL({MdpJsonSql.Int("m", "RecID")},0)),
+                @BatchId, @Now
+            FROM mdp_stg_production_return m
+            LEFT JOIN DepartmentMaster d
+              ON d.Domain = {MdpJsonSql.Str("m", "Domain")} AND d.Department = {MdpJsonSql.Str("m", "Department")}
+            WHERE {headWhere}
             ON DUPLICATE KEY UPDATE
                 factory_id=VALUES(factory_id), nbr=VALUES(nbr), return_date=VALUES(return_date),
                 status=VALUES(status), status_desc=VALUES(status_desc), work_ord=VALUES(work_ord),
@@ -174,45 +366,65 @@ public class ProductionReturnMdpSyncService : ITransient
                 source_create_time=VALUES(source_create_time), trans_type=VALUES(trans_type),
                 trans_type_text=VALUES(trans_type_text), eff_date=VALUES(eff_date), ufld1=VALUES(ufld1),
                 sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
-            """,
-            new SugarParameter("@BatchId", batchId),
-            new SugarParameter("@Now", now));
+            """;
+        var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
+        if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
+        await _db.Ado.ExecuteCommandAsync(insertSql, insPars);
         return rows;
     }
 
     /// <summary>
-    /// 标准化明细:NbrDetail(WOD) + ItemMaster + LocationMaster -> mdp_std_production_return_detail(回填 std_head_id)。
-    /// 返回处理明细行数。LocationMaster 列名小写 location/descr,分别 join LocationFrom/LocationTo
+    /// 标准化明细:stg(NbrDetail, WOD) self-join stg(NbrMaster) by RecID↔NbrRecID + ItemMaster/LocationMaster(本库)
+    /// → mdp_std_production_return_detail(回填 std_head_id)。返回处理明细行数
     /// </summary>
-    private async Task<int> TransformDetailStandardAsync(string batchId, DateTime now)
+    private async Task<int> TransformDetailStandardAsync(string batchId, DateTime now, string? sourceSystem = null)
     {
+        var srcClause = sourceSystem == null ? "" : " AND m.source_system=@Src AND d.source_system=@Src";
+        var dTenant = $"IFNULL({MdpJsonSql.Int("d", "tenant_id")}, IFNULL(d.tenant_id, 0))";
+        var dStatusExpr = $"UPPER(IFNULL({MdpJsonSql.Str("d", "Status")}, ''))";
+        // 明细 from/join/where 片段(count 与 insert 共用,避免发散)。
+        var fromJoinWhere =
+            $"""
+            FROM mdp_stg_production_return d
+            JOIN mdp_stg_production_return m
+              ON m.source_table='NbrMaster'
+             AND {MdpJsonSql.Int("m", "RecID")} = {MdpJsonSql.Int("d", "NbrRecID")}
+             AND {MdpJsonSql.Str("m", "Type")}='WOD'
+             AND {MdpJsonSql.BoolTrue("m", "IsActive")}
+             AND LOWER({MdpJsonSql.Raw("m", "IsReturn")}) IN ('0','false')
+            LEFT JOIN ItemMaster i
+              ON i.Domain = {MdpJsonSql.Str("d", "Domain")} AND i.ItemNum = {MdpJsonSql.Str("d", "ItemNum")}
+            LEFT JOIN LocationMaster lf
+              ON lf.tenant_id = {dTenant} AND lf.location = {MdpJsonSql.Str("d", "LocationFrom")}
+            LEFT JOIN LocationMaster lt
+              ON lt.tenant_id = {dTenant} AND lt.location = {MdpJsonSql.Str("d", "LocationTo")}
+            LEFT JOIN mdp_std_production_return h
+              ON h.tenant_id = {dTenant} AND h.domain = {MdpJsonSql.Str("m", "Domain")} AND h.rec_id = {MdpJsonSql.Int("d", "NbrRecID")}
+            WHERE d.source_table='NbrDetail' AND {MdpJsonSql.Str("d", "Type")}='WOD'{srcClause}
+            """;
+
+        var countPars = new List<SugarParameter>();
+        if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
         var rows = await _db.Ado.GetIntAsync(
-            """
-            SELECT COUNT(1) FROM NbrDetail dt
-            JOIN NbrMaster m ON m.RecID = dt.NbrRecID AND m.Type='WOD' AND m.IsActive=1 AND m.IsReturn=0
-            WHERE dt.Type='WOD'
-            """);
-        await _db.Ado.ExecuteCommandAsync(
-            """
+            $"SELECT COUNT(1) {fromJoinWhere}", countPars);
+
+        var insertSql =
+            $"""
             INSERT INTO mdp_std_production_return_detail
             (tenant_id, std_head_id, domain, nbr_rec_id, rec_id, line, item_num, item_name, item_spec, um,
              qty_to, qty_rec, location_from, location_from_desc, location_to, location_to_desc,
              lot_serial, status, status_desc, remark, work_ord, source_biz_key, sync_batch_id, sync_time)
             SELECT
-                IFNULL(dt.tenant_id, 0), h.id, m.Domain, dt.NbrRecID, dt.RecID, dt.Line,
-                dt.ItemNum, i.Descr, i.Descr1, dt.UM,
-                dt.QtyTo, dt.QtyRec, dt.LocationFrom, lf.descr, dt.LocationTo, lt.descr,
-                CAST(dt.LotSerial AS CHAR), UPPER(IFNULL(dt.Status, '')),
-                CASE WHEN UPPER(IFNULL(dt.Status,''))='C' THEN '关闭' ELSE UPPER(IFNULL(dt.Status,'')) END,
-                CAST(dt.Remark AS CHAR), dt.WorkOrd,
-                CONCAT(m.Domain, ':', dt.NbrRecID, ':', dt.Line), @BatchId, @Now
-            FROM NbrDetail dt
-            JOIN NbrMaster m ON m.RecID = dt.NbrRecID AND m.Type='WOD' AND m.IsActive=1 AND m.IsReturn=0
-            LEFT JOIN ItemMaster i ON i.Domain = dt.Domain AND i.ItemNum = dt.ItemNum
-            LEFT JOIN LocationMaster lf ON lf.tenant_id = IFNULL(dt.tenant_id, 0) AND lf.location = dt.LocationFrom
-            LEFT JOIN LocationMaster lt ON lt.tenant_id = IFNULL(dt.tenant_id, 0) AND lt.location = dt.LocationTo
-            LEFT JOIN mdp_std_production_return h ON h.tenant_id = IFNULL(dt.tenant_id, 0) AND h.domain = m.Domain AND h.rec_id = dt.NbrRecID
-            WHERE dt.Type='WOD'
+                {dTenant}, h.id, {MdpJsonSql.Str("m", "Domain")}, {MdpJsonSql.Int("d", "NbrRecID")}, {MdpJsonSql.Int("d", "RecID")}, {MdpJsonSql.Int("d", "Line")},
+                {MdpJsonSql.Str("d", "ItemNum")}, i.Descr, i.Descr1, {MdpJsonSql.Str("d", "UM")},
+                {MdpJsonSql.Dec("d", "QtyTo", 18, 5)}, {MdpJsonSql.Dec("d", "QtyRec", 18, 5)},
+                {MdpJsonSql.Str("d", "LocationFrom")}, lf.descr, {MdpJsonSql.Str("d", "LocationTo")}, lt.descr,
+                {MdpJsonSql.Str("d", "LotSerial")}, {dStatusExpr},
+                CASE WHEN {dStatusExpr}='C' THEN '关闭' ELSE {dStatusExpr} END,
+                {MdpJsonSql.Str("d", "Remark")}, {MdpJsonSql.Str("d", "WorkOrd")},
+                CONCAT(IFNULL({MdpJsonSql.Str("m", "Domain")},''), ':', IFNULL({MdpJsonSql.Int("d", "NbrRecID")},0), ':', IFNULL({MdpJsonSql.Int("d", "Line")},0)),
+                @BatchId, @Now
+            {fromJoinWhere}
             ON DUPLICATE KEY UPDATE
                 std_head_id=VALUES(std_head_id), rec_id=VALUES(rec_id), item_num=VALUES(item_num),
                 item_name=VALUES(item_name), item_spec=VALUES(item_spec), um=VALUES(um),
@@ -221,9 +433,10 @@ public class ProductionReturnMdpSyncService : ITransient
                 location_to_desc=VALUES(location_to_desc), lot_serial=VALUES(lot_serial),
                 status=VALUES(status), status_desc=VALUES(status_desc), remark=VALUES(remark), work_ord=VALUES(work_ord),
                 sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
-            """,
-            new SugarParameter("@BatchId", batchId),
-            new SugarParameter("@Now", now));
+            """;
+        var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
+        if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
+        await _db.Ado.ExecuteCommandAsync(insertSql, insPars);
         return rows;
     }
 
@@ -288,3 +501,15 @@ public sealed class ProductionReturnMdpSyncResult
     public int HeadRows { get; set; }
     public int DetailRows { get; set; }
 }
+
+/// <summary>S5 生产退料双模式入站结果。</summary>
+public sealed class ProductionReturnInboundResult
+{
+    public string PullBatchId { get; set; } = string.Empty;
+    public int RowsPulled { get; set; }
+    public int RowsWrittenStg { get; set; }
+    public string? NewCursor { get; set; }
+    public string TransformBatchId { get; set; } = string.Empty;
+    public int StdHeadRows { get; set; }
+    public int StdDetailRows { get; set; }
+}