using Admin.NET.Plugin.AiDOP.DataPlatform.Executors; namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse; /// /// S5 生产入库单 数据中台只读同步转换服务(DOP 内部,扁平单表,独立于 S5MdpSyncTransformService 的 KPI 管线)。 /// /// 源:aidopdev NbrMaster(n) + NbrDetail(d),业务类型 Type='WOI'(生产入库); /// 扁平 LIST join:n.Domain=d.Domain AND n.Nbr=d.Nbr(一行=一条 NbrDetail 明细)。 /// 维表 DepartmentMaster(p) / ItemMaster(i) / LocationMaster(lt,lf) 全 LEFT JOIN。 /// LocationMaster 无 Domain 列(DOP 重建:domain_code/location/descr),按 tenant_id+location join。 /// 链路(Phase 1):执行器抽 NbrMaster/NbrDetail → mdp_stg_production_receipt → mdp_std_production_receipt。 /// 双模式入站:mdp_entity=S7_NBR_MASTER/S7_NBR_DETAIL → Writer → stg;标准层读 stg;维表仍直连本库。 /// /// 约束: /// - 只读源/贴源,仅写 mdp_stg_production_receipt / mdp_std_production_receipt;绝不写 NbrMaster/NbrDetail。 /// - 过滤 Type='WOI' AND IsActive=1。 /// - WOI 为 0 行时转换成功完成、处理数为 0,不报错。 /// public class ProductionReceiptMdpSyncService : ITransient { private const string JobCode = "S5_PRODUCTION_RECEIPT_MDP_SYNC"; private const string InboundEntityCode = "S7_NBR_MASTER"; private const string InboundDetailEntityCode = "S7_NBR_DETAIL"; private readonly ISqlSugarClient _db; private readonly MdpSourcePullDispatcher _pullDispatcher; public ProductionReceiptMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher) { _db = db; _pullDispatcher = pullDispatcher; } /// 全量:本地 DB 执行器灌 stg → 标准层(读 stg)。 public async Task RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO") { cancellationToken.ThrowIfCancellationRequested(); await EnsureTablesAsync(); await EnsureStgTableAsync(); var now = DateTime.Now; var batchId = $"S5_PROD_RCPT_FULL_{now:yyyyMMddHHmmss}"; var runLogId = await InsertRunLogAsync(batchId, now, triggerType); var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId }; try { var pullCtx = new MdpPullContext { TenantId = 0, FullRefresh = true, TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND", BatchId = $"{batchId}_PULL" }; await PopulateStgAsync(pullCtx, cancellationToken); result.Rows = await TransformStandardAsync(batchId, now); await MarkRunSuccessAsync(runLogId, now, result); return result; } catch (Exception ex) { await MarkRunFailedAsync(runLogId, now, ex.Message); throw; } } /// /// S7 成品入库双模式入站:执行器抽主/明细 → stg,再跑 WOI 标准层(读 stg)。 /// public async Task 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 = "S7_PRODUCTION_RECEIPT_INBOUND", BatchId = $"S7_NBR_IN_{now:yyyyMMddHHmmss}" }; var pull = await PopulateStgAsync(pullCtx, cancellationToken); var batchId = $"S7_NBR_STD_{now:yyyyMMddHHmmss}"; var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND"); var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId }; try { result.Rows = await TransformStandardAsync(batchId, now); await MarkRunSuccessAsync(runLogId, now, result); } catch (Exception ex) { await MarkRunFailedAsync(runLogId, now, ex.Message); throw; } return new ProductionReceiptInboundResult { PullBatchId = pullCtx.BatchId, RowsPulled = pull.RowsPulled, RowsWrittenStg = pull.RowsWritten, NewCursor = pull.NewCursor, TransformBatchId = result.BatchId, StdRows = result.Rows }; } private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync( MdpPullContext pullCtx, CancellationToken cancellationToken) { var master = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, 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_receipt ( 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='S7成品入库执行器贴源层' """); } /// 防御式建表(与 UpdateScripts/1.0.214.sql 同构,幂等)。 private async Task EnsureTablesAsync() { await _db.Ado.ExecuteCommandAsync( """ CREATE TABLE IF NOT EXISTS mdp_std_production_receipt ( id BIGINT AUTO_INCREMENT PRIMARY KEY, tenant_id BIGINT NOT NULL DEFAULT 0, factory_id BIGINT NULL DEFAULT 1, source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP', domain VARCHAR(80) NOT NULL, master_rec_id INT NULL, detail_rec_id INT NOT NULL, nbr VARCHAR(24) NULL, line SMALLINT NOT NULL DEFAULT 0, receipt_date DATETIME NULL, status VARCHAR(8) NULL, status_desc VARCHAR(20) NULL, remark VARCHAR(200) NULL, prod_line VARCHAR(8) NULL, work_ord VARCHAR(64) NULL, erp_work_ord VARCHAR(60) NULL, department VARCHAR(8) NULL, department_desc VARCHAR(255) NULL, applicant_name VARCHAR(12) NULL, item_num VARCHAR(24) NULL, item_name VARCHAR(200) NULL, item_spec VARCHAR(200) NULL, um VARCHAR(8) NULL, location_to VARCHAR(8) NULL, location_to_desc VARCHAR(255) NULL, lot_serial VARCHAR(120) NULL, qty_rec DECIMAL(18,5) NULL DEFAULT 0, qty_to DECIMAL(18,5) NULL DEFAULT 0, location_from VARCHAR(8) NULL, location_from_desc VARCHAR(255) NULL, ord_nbr VARCHAR(48) NULL, source_biz_key VARCHAR(200) NULL, sync_batch_id VARCHAR(100) NOT NULL, sync_time DATETIME NOT NULL, create_time DATETIME DEFAULT CURRENT_TIMESTAMP, update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_mdp_std_prod_receipt (tenant_id, domain, detail_rec_id), KEY idx_mdp_std_prod_rcpt_nbr (tenant_id, nbr), KEY idx_mdp_std_prod_rcpt_date (tenant_id, receipt_date), KEY idx_mdp_std_prod_rcpt_workord (tenant_id, work_ord), KEY idx_mdp_std_prod_rcpt_erp (tenant_id, erp_work_ord), KEY idx_mdp_std_prod_rcpt_lot (tenant_id, lot_serial), KEY idx_mdp_std_prod_rcpt_item (tenant_id, item_num), KEY idx_mdp_std_prod_rcpt_locto (tenant_id, location_to) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S5生产入库单标准层(扁平,一行=一条明细)' """); } /// /// 标准化:stg(NbrMaster/NbrDetail) + 维表(本库) → mdp_std_production_receipt。 /// private async Task TransformStandardAsync(string batchId, DateTime now) { var rows = await _db.Ado.GetIntAsync( """ SELECT COUNT(1) FROM mdp_stg_production_receipt n INNER JOIN mdp_stg_production_receipt d ON d.source_table='NbrDetail' AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr')) WHERE n.source_table='NbrMaster' AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI' AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1 AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL """); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_production_receipt (tenant_id, factory_id, source_system, domain, master_rec_id, detail_rec_id, nbr, line, receipt_date, status, status_desc, remark, prod_line, work_ord, erp_work_ord, department, department_desc, applicant_name, item_num, item_name, item_spec, um, location_to, location_to_desc, lot_serial, qty_rec, qty_to, location_from, location_from_desc, ord_nbr, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(n.tenant_id, 0), 1, IFNULL(NULLIF(n.source_system,''), 'AIDOP'), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.RecID')) AS SIGNED), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS SIGNED), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS SIGNED), STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Date')),'null'),''), '%Y-%m-%d %H:%i:%s'), UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')), UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Remark')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ProdLine')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.WorkOrd')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Address')), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')), TRIM(CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),''), ' ', IFNULL(p.Descr, ''))), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Name')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')), i.Descr, i.Descr1, i.UM, JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')), lt.descr, CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LotSerial')) AS CHAR), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyRec')) AS DECIMAL(18,5)), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyTo')) AS DECIMAL(18,5)), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')), lf.descr, JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')), CONCAT(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), '#', JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID'))), @BatchId, @Now FROM mdp_stg_production_receipt n INNER JOIN mdp_stg_production_receipt d ON d.source_table='NbrDetail' AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr')) LEFT JOIN DepartmentMaster p ON p.Domain = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) AND p.Department = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')) LEFT JOIN ItemMaster i ON i.Domain = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) AND i.ItemNum = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')) LEFT JOIN LocationMaster lt ON lt.tenant_id = IFNULL(d.tenant_id, 0) AND lt.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')) LEFT JOIN LocationMaster lf ON lf.tenant_id = IFNULL(d.tenant_id, 0) AND lf.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')) WHERE n.source_table='NbrMaster' AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI' AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1 AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL ON DUPLICATE KEY UPDATE factory_id=VALUES(factory_id), master_rec_id=VALUES(master_rec_id), nbr=VALUES(nbr), line=VALUES(line), receipt_date=VALUES(receipt_date), status=VALUES(status), status_desc=VALUES(status_desc), remark=VALUES(remark), prod_line=VALUES(prod_line), work_ord=VALUES(work_ord), erp_work_ord=VALUES(erp_work_ord), department=VALUES(department), department_desc=VALUES(department_desc), applicant_name=VALUES(applicant_name), item_num=VALUES(item_num), item_name=VALUES(item_name), item_spec=VALUES(item_spec), um=VALUES(um), location_to=VALUES(location_to), location_to_desc=VALUES(location_to_desc), lot_serial=VALUES(lot_serial), qty_rec=VALUES(qty_rec), qty_to=VALUES(qty_to), location_from=VALUES(location_from), location_from_desc=VALUES(location_from_desc), ord_nbr=VALUES(ord_nbr), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP """, new SugarParameter("@BatchId", batchId), new SugarParameter("@Now", now)); return rows; } private async Task InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType) { await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_transform_run_log (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time) VALUES (0, @JobCode, 'S5生产入库单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime) """, new SugarParameter("@JobCode", JobCode), new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)), new SugarParameter("@BatchId", batchId), new SugarParameter("@StartTime", startedAt)); return await _db.Ado.GetLongAsync( "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1", new List { new("@BatchId", batchId) }); } private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, ProductionReceiptMdpSyncResult result) { var finishedAt = DateTime.Now; await _db.Ado.ExecuteCommandAsync( """ UPDATE mdp_transform_run_log SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs, stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP WHERE id=@Id """, new SugarParameter("@EndTime", finishedAt), new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds), new SugarParameter("@StandardRows", result.Rows), new SugarParameter("@Id", runLogId)); } private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message) { var finishedAt = DateTime.Now; await _db.Ado.ExecuteCommandAsync( """ UPDATE mdp_transform_run_log SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs, error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP WHERE id=@Id """, new SugarParameter("@EndTime", finishedAt), new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds), new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message), new SugarParameter("@Id", runLogId)); } private static string NormalizeTriggerType(string? triggerType) => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant(); } /// 生产入库单 MDP 同步转换结果。 public sealed class ProductionReceiptMdpSyncResult { public long RunLogId { get; set; } public string BatchId { get; set; } = string.Empty; public int Rows { get; set; } } /// S7 成品入库双模式入站结果。 public sealed class ProductionReceiptInboundResult { 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 StdRows { get; set; } }