using System.Text.Json; using Admin.NET.Plugin.AiDOP.DataPlatform; using Admin.NET.Plugin.AiDOP.DataPlatform.Executors; using Admin.NET.Plugin.AiDOP.Infrastructure; namespace Admin.NET.Plugin.AiDOP.Manufacturing; /// /// S6 报工双模式入站:执行器 → mdp_stg_s6_report → mdp_std_s6_report。 /// 源优先 T8 Cj_Bg_Head_Rep(实体 S6_REPORT);API 对偶 S6_REPORT_API。 /// public class ReportWorkMdpSyncService : ITransient { private const string InboundEntityCode = "S6_REPORT"; /// T8 报工源无 tenant_id;须经 解析。 private const string SourceZtid = "pbxfxp"; private readonly ISqlSugarClient _db; private readonly MdpSourcePullDispatcher _pullDispatcher; private readonly MdpNeutralSourceGate _neutralGate; public ReportWorkMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher, MdpNeutralSourceGate neutralGate) { _db = db; _pullDispatcher = pullDispatcher; _neutralGate = neutralGate; } public async Task RunInboundAsync( long tenantId = 0, bool fullRefresh = false, string? entityCode = null, CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); await EnsureTablesAsync(); if (tenantId > 0 && await _neutralGate.AllowsAsync(tenantId, "S6_REPORT", "DOPDEMORQ_SQLSERVER")) return await Run165Async(tenantId, fullRefresh, cancellationToken); // T8 源表无 tenant_id 列;显式 tenantId≤0 时经账套映射解析真实租户。 var resolvedTenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId); var code = string.IsNullOrWhiteSpace(entityCode) ? InboundEntityCode : entityCode.Trim(); var now = DateTime.Now; var pullCtx = new MdpPullContext { TenantId = resolvedTenantId, FullRefresh = fullRefresh, TaskCode = "S6_REPORT_INBOUND", BatchId = $"S6_RPT_IN_{now:yyyyMMddHHmmss}" }; var pull = await _pullDispatcher.PullByEntityCodeAsync(code, pullCtx, cancellationToken); var transformBatch = $"{pullCtx.BatchId}_STD"; var stdRows = await TransformStandardAsync(resolvedTenantId, transformBatch, now); return new ReportWorkInboundResult { PullBatchId = pullCtx.BatchId, RowsPulled = pull.RowsWritten, RowsWrittenStg = pull.RowsWritten, TransformBatchId = transformBatch, StandardRows = stdRows }; } private async Task TransformStandardAsync(long tenantId, string batchId, DateTime now) { if (!await _neutralGate.AllowsAsync(tenantId, "S6_REPORT", "T8")) return 0; // 将 PENDING stg 投影到 std(JSON 字段兼容大小写) var sTenant = MdpJsonSql.TenantFromStg("s", "@tid"); var sql = $""" INSERT INTO mdp_std_s6_report (tenant_id, factory_id, source_system, work_order_no, report_date, report_qty, ztid, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT {sTenant}, 1, IFNULL(NULLIF(s.source_system,''), 'T8'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.noid')), 'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.NOID')), 'null'), s.source_biz_key)), COALESCE( STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.kgdate')), 'null'), ''), '%Y-%m-%d %H:%i:%s'), STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.KGDATE')), 'null'), ''), '%Y-%m-%d %H:%i:%s'), NULL), CAST(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.sl')), 'null'), '') AS DECIMAL(18,6)), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data, '$.ztid')), 'null'), IFNULL(NULLIF(s.source_row_id,''), s.source_biz_key), s.source_biz_key, @batch, @now FROM mdp_stg_s6_report s WHERE s.process_status = 'PENDING' AND s.source_biz_key IS NOT NULL AND s.source_biz_key <> '' AND {MdpJsonSql.TenantGuard(sTenant)} ON DUPLICATE KEY UPDATE work_order_no = VALUES(work_order_no), report_date = VALUES(report_date), report_qty = VALUES(report_qty), ztid = VALUES(ztid), source_row_id = VALUES(source_row_id), sync_batch_id = VALUES(sync_batch_id), sync_time = VALUES(sync_time), update_time = CURRENT_TIMESTAMP """; var affected = await MdpSchemaAligner.ExecuteAsync(_db, sql, new SugarParameter("@tid", tenantId), new SugarParameter("@batch", batchId), new SugarParameter("@now", now)); await MdpSchemaAligner.ExecuteAsync(_db, tenantId > 0 ? "UPDATE mdp_stg_s6_report SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING' AND tenant_id = @tid" : "UPDATE mdp_stg_s6_report SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING'", tenantId > 0 ? new SugarParameter("@tid", tenantId) : null!); return affected; } private async Task Run165Async(long tenantId, bool fullRefresh, CancellationToken cancellationToken) { var now = DateTime.Now; var pullCtx = new MdpPullContext { TenantId = tenantId, FullRefresh = fullRefresh, TaskCode = "S6_EMP_WORK_ORD_TRANS", BatchId = $"S6_165_RPT_{now:yyyyMMddHHmmss}" }; var pull = await _pullDispatcher.PullByEntityCodeAsync("S6_EMP_WORK_ORD_TRANS_SQLSERVER", pullCtx, cancellationToken); await EnsureStartWorkColumnAsync(); var batch = $"{pullCtx.BatchId}_STD"; var stdRows = await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_s6_report (tenant_id, factory_id, source_system, written_by, work_order_no, report_date, start_work_date, report_qty, ztid, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, 'DOPDEMORQ_SQLSERVER', 'DB_SYNC', NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrd')),'null'), COALESCE( STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), STR_TO_DATE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdDate')),'null'), '%Y-%m-%d')), COALESCE( STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.StartDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), STR_TO_DATE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.StartDate')),'null'), '%Y-%m-%d')), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CompQty')),'null') AS DECIMAL(18,6)), '', IFNULL(NULLIF(s.source_row_id,''), s.source_biz_key), s.source_biz_key, @batch, @now FROM mdp_stg_s6_report s WHERE s.tenant_id=@tid AND s.process_status='PENDING' AND s.source_system='DOPDEMORQ_SQLSERVER' ON DUPLICATE KEY UPDATE report_qty=VALUES(report_qty), report_date=VALUES(report_date), start_work_date=VALUES(start_work_date), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@batch", batch), new SugarParameter("@now", now)); await _db.Ado.ExecuteCommandAsync( """ UPDATE mdp_std_work_order_line w JOIN ( SELECT work_order_no, MIN(start_work_date) d FROM mdp_std_s6_report WHERE tenant_id=@tid AND source_system='DOPDEMORQ_SQLSERVER' AND start_work_date IS NOT NULL GROUP BY work_order_no ) r ON r.work_order_no=w.order_no SET w.release_time=r.d WHERE w.tenant_id=@tid AND w.doc_type='PROD_TASK' AND w.source_system='AIDOP_NATIVE' """, new SugarParameter("@tid", tenantId)); await _db.Ado.ExecuteCommandAsync( "UPDATE mdp_stg_s6_report SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING' AND tenant_id=@tid AND source_system='DOPDEMORQ_SQLSERVER'", new SugarParameter("@tid", tenantId)); return new ReportWorkInboundResult { PullBatchId = pullCtx.BatchId, RowsPulled = pull.RowsWritten, RowsWrittenStg = pull.RowsWritten, TransformBatchId = batch, StandardRows = stdRows }; } private async Task EnsureStartWorkColumnAsync() { var exists = await _db.Ado.GetIntAsync( "SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='mdp_std_s6_report' AND COLUMN_NAME='start_work_date'"); if (exists == 0) await _db.Ado.ExecuteCommandAsync("ALTER TABLE mdp_std_s6_report ADD COLUMN start_work_date datetime NULL"); await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_s6_report"); } private async Task EnsureTablesAsync() { await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_s6_report")); await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_s6_report")); } } public class ReportWorkInboundResult { public string PullBatchId { get; set; } = string.Empty; public int RowsPulled { get; set; } public int RowsWrittenStg { get; set; } public string TransformBatchId { get; set; } = string.Empty; public int StandardRows { get; set; } }