using Admin.NET.Plugin.AiDOP.DataPlatform.Schema; using Admin.NET.Plugin.AiDOP.Infrastructure; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.DataPlatform; /// /// 自有业务表投影到中立层。只写登记为 AIDOP_NATIVE 的对象。 /// public sealed class NativeNeutralProjectionService : ITransient { private readonly ISqlSugarClient _db; private readonly MdpNeutralSourceGate _neutralGate; private readonly EmployeePositionMapService _positions; private readonly SourceDomainTenantResolver _domains; private readonly ILogger _logger; public NativeNeutralProjectionService( ISqlSugarClient db, MdpNeutralSourceGate neutralGate, EmployeePositionMapService positions, SourceDomainTenantResolver domains, ILoggerFactory loggerFactory) { _db = db; _neutralGate = neutralGate; _positions = positions; _domains = domains; _logger = loggerFactory.CreateLogger(nameof(NativeNeutralProjectionService)); } /// /// 销售发货实绩(登记来源的 SALES_SHIP)按「订单号 + 物料」汇总。 /// /// 模式一下销售订单的权威在自有 S1,执行事实由外部系统提供,故订单行的已交货量与关行状态 /// 要用发货流水推导,不能只看 mdp_std_so.delivered_qty(自建单常年为空)。 /// 来源不写死:内联 mdp_tenant_std_source,只读该租户登记为 INV_TRANS 权威的那一个来源。 /// 粒度与 S7 指标 SQL 的 JOIN 一致(中立流水没有订单行号,最细只能到订单号 + 物料)。 /// /// private const string ShipmentActualsSql = """ SELECT t.tenant_id, t.ref_task_no, t.item_num, SUM(ABS(IFNULL(t.qty_change,0))) AS shipped_qty, MAX(t.approved_time) AS last_ship_time FROM mdp_std_inv_trans t INNER JOIN mdp_tenant_std_source ts ON ts.tenant_id=t.tenant_id AND ts.std_object='INV_TRANS' AND ts.source_system=t.source_system WHERE t.tenant_id=@tid AND t.biz_doc_type='SALES_SHIP' AND t.summary_flag=0 AND t.void_flag=0 AND t.approved_flag=1 GROUP BY t.tenant_id, t.ref_task_no, t.item_num """; public async Task ProjectAsync(long tenantId, string batchId, CancellationToken cancellationToken = default) { if (tenantId <= 0) return; var failed = new List(); await Try(failed, "WO_LINE_PROD", () => ProjectWorkOrdersAsync(tenantId, batchId, cancellationToken)); await Try(failed, "WO_SCHEDULE", () => ProjectWorkOrderScheduleAsync(tenantId, batchId, cancellationToken)); await Try(failed, "WO_BOM", () => ProjectBomAsync(tenantId, batchId, cancellationToken)); await Try(failed, "EMPLOYEE", () => ProjectEmployeesAsync(tenantId, batchId, cancellationToken)); await Try(failed, "WO_LINE_SALES", () => ProjectSalesLinesAsync(tenantId, batchId, cancellationToken)); if (failed.Count > 0) _logger.LogWarning("自有中立投影未完成 tenant={Tenant} {Failed}", tenantId, string.Join(" | ", failed)); } private async Task Try(List failed, string name, Func action) { try { await action(); } catch (Exception ex) when (ex is not OperationCanceledException) { failed.Add($"{name}: {ex.Message}"); } } private async Task ProjectWorkOrdersAsync(long tenantId, string batchId, CancellationToken cancellationToken) { if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_PROD", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken)) return; await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_line"); cancellationToken.ThrowIfCancellationRequested(); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_work_order_line (tenant_id, factory_id, source_system, written_by, domain, doc_type, src_doc_type_raw, order_no, line_no, task_no, item_code, qty_planned, qty_completed, plan_finish_date, release_time, closed_flag, closed_time, approved_flag, void_flag, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT w.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(w.Domain,''), 'PROD_TASK', 'WorkOrdMaster', w.WorkOrd, NULL, w.WorkOrd, w.ItemNum, w.QtyOrded, w.QtyCompleted, w.DueDate, w.ReleaseDate, IF(UPPER(IFNULL(w.Status,''))='C', 1, 0), NULL, 1, IF(IFNULL(w.IsActive,1)=0, 1, 0), CAST(w.RecID AS CHAR), LEFT(CONCAT(IFNULL(w.Domain,''), ':', w.WorkOrd), 200), @batch, NOW() FROM WorkOrdMaster w WHERE w.tenant_id=@tid AND w.WorkOrd IS NOT NULL AND w.WorkOrd<>'' ON DUPLICATE KEY UPDATE item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), qty_completed=VALUES(qty_completed), plan_finish_date=VALUES(plan_finish_date), release_time=VALUES(release_time), closed_flag=VALUES(closed_flag), void_flag=VALUES(void_flag), written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time) """, new { tid = tenantId, batch = batchId }); } /// /// 自有生产工单 → mdp_std_work_order_schedule(doc_type='PROD_TASK')。 /// /// S5_L1_002 物料齐套满足率的分母以工单排程头为驱动,且 BOM 按 source_system 关联, /// 故排程必须与 写同一个 AIDOP_NATIVE,否则分母恒为空、 /// 指标只会算出 NO_DATA。 /// /// /// 唯一键 uk_std_wo_sched_type 不含 source_system,同一工单号若已被别的来源占用, /// 直接 upsert 会把对方的行改成自有来源,故与 T8 投影一样加 NOT EXISTS 反向占用保护。 /// /// private async Task ProjectWorkOrderScheduleAsync(long tenantId, string batchId, CancellationToken cancellationToken) { if (!await _neutralGate.AllowsAsync(tenantId, "WO_SCHEDULE", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken)) return; await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_schedule"); cancellationToken.ThrowIfCancellationRequested(); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_work_order_schedule (tenant_id, factory_id, source_system, written_by, work_order, doc_type, item_code, item_name, site_code, status, priority, urgent_flag, qty_ordered, qty_completed, order_date, due_date, release_date, prod_line, lot_serial, drawing_no, project, work_order_type, labor_variance, approved_flag, void_flag, source_biz_key, sync_batch_id, sync_time) SELECT w.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', w.WorkOrd, 'PROD_TASK', w.ItemNum, w.ItemName, w.Site, w.Status, w.Priority, IF(IFNULL(w.Urgent,0)<>0, 1, 0), w.QtyOrded, w.QtyCompleted, w.OrdDate, w.DueDate, w.ReleaseDate, w.ProdLine, w.Batch, w.Drawing, w.Project, w.Typed, w.LbrVar, 1, IF(IFNULL(w.IsActive,1)=0, 1, 0), LEFT(CONCAT(IFNULL(w.Domain,''), ':', w.WorkOrd), 200), @batch, NOW() FROM WorkOrdMaster w WHERE w.tenant_id=@tid AND w.WorkOrd IS NOT NULL AND w.WorkOrd<>'' AND NOT EXISTS ( SELECT 1 FROM mdp_std_work_order_schedule x WHERE x.tenant_id=w.tenant_id AND x.doc_type='PROD_TASK' AND x.work_order=w.WorkOrd AND x.source_system<>'AIDOP_NATIVE') ON DUPLICATE KEY UPDATE item_code=VALUES(item_code), item_name=VALUES(item_name), site_code=VALUES(site_code), status=VALUES(status), priority=VALUES(priority), urgent_flag=VALUES(urgent_flag), qty_ordered=VALUES(qty_ordered), qty_completed=VALUES(qty_completed), order_date=VALUES(order_date), due_date=VALUES(due_date), release_date=VALUES(release_date), prod_line=VALUES(prod_line), lot_serial=VALUES(lot_serial), drawing_no=VALUES(drawing_no), project=VALUES(project), work_order_type=VALUES(work_order_type), labor_variance=VALUES(labor_variance), approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag), written_by=VALUES(written_by), source_biz_key=VALUES(source_biz_key), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time) """, new { tid = tenantId, batch = batchId }); } private async Task ProjectBomAsync(long tenantId, string batchId, CancellationToken cancellationToken) { if (!await _neutralGate.AllowsAsync(tenantId, "WO_BOM", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken)) return; await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_bom"); cancellationToken.ThrowIfCancellationRequested(); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_work_order_bom (tenant_id, factory_id, source_system, written_by, domain, order_no, order_line_no, task_no, item_code, qty_required, unit, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT d.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(d.Domain,''), d.WorkOrd, CAST(d.LineNum AS CHAR), d.WorkOrd, d.ItemNum, d.QtyRequired, d.UM, CAST(d.RecID AS CHAR), LEFT(CONCAT(IFNULL(d.Domain,''), ':', d.WorkOrd, ':', d.LineNum, ':', IFNULL(d.ItemNum,'')), 200), @batch, NOW() FROM WorkOrdDetail d WHERE d.tenant_id=@tid AND d.WorkOrd IS NOT NULL AND d.WorkOrd<>'' ON DUPLICATE KEY UPDATE item_code=VALUES(item_code), qty_required=VALUES(qty_required), unit=VALUES(unit), written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time) """, new { tid = tenantId, batch = batchId }); } private async Task ProjectEmployeesAsync(long tenantId, string batchId, CancellationToken cancellationToken) { if (!await _neutralGate.AllowsAsync(tenantId, "EMPLOYEE", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken)) return; await _positions.SyncNativeAsync(tenantId, cancellationToken); await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_employee"); cancellationToken.ThrowIfCancellationRequested(); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_employee (tenant_id, factory_id, source_system, written_by, domain, employee_no, employee_name, department_code, position_code, src_position_raw, employment_status, src_employment_status_raw, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT e.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(e.Domain,''), e.Employee, e.Name, e.Department, IFNULL(m.position_code, 'UNKNOWN'), COALESCE(NULLIF(e.JobTitle,''), NULLIF(e.WorkCtr,'')), CASE WHEN e.DateTerminated IS NOT NULL AND e.DateTerminated <= CURDATE() THEN 'LEFT' WHEN e.EmploymentStatus IN ('在职') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('ACTIVE','ONJOB','ON_JOB') THEN 'ACTIVE' WHEN e.EmploymentStatus IN ('离职','辞退') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('LEFT','LEAVE','RESIGNED','TERMINATED','QUIT') THEN 'LEFT' WHEN e.EmploymentStatus IN ('停用') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('INACTIVE','DISABLED','SUSPENDED') THEN 'INACTIVE' ELSE 'UNKNOWN' END, e.EmploymentStatus, CAST(e.RecID AS CHAR), LEFT(CONCAT(IFNULL(e.Domain,''), ':', e.Employee), 200), @batch, NOW() FROM EmployeeMaster e LEFT JOIN mdp_employee_position_map m ON m.tenant_id=e.tenant_id AND m.source_system='AIDOP_NATIVE' AND m.domain=IFNULL(e.Domain,'') AND m.src_position_raw=COALESCE(NULLIF(e.JobTitle,''), NULLIF(e.WorkCtr,'')) WHERE e.tenant_id=@tid AND e.Employee IS NOT NULL AND e.Employee<>'' ON DUPLICATE KEY UPDATE employee_name=VALUES(employee_name), department_code=VALUES(department_code), position_code=VALUES(position_code), src_position_raw=VALUES(src_position_raw), employment_status=VALUES(employment_status), src_employment_status_raw=VALUES(src_employment_status_raw), written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time) """, new { tid = tenantId, batch = batchId }); } private async Task ProjectSalesLinesAsync(long tenantId, string batchId, CancellationToken cancellationToken) { if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_SALES", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken)) return; var domain = await _domains.ResolveDomainAsync(MdpSourceIdentity.Native, tenantId, cancellationToken); await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_line"); cancellationToken.ThrowIfCancellationRequested(); await _db.Ado.ExecuteCommandAsync( $""" INSERT INTO mdp_std_work_order_line (tenant_id, factory_id, source_system, written_by, domain, doc_type, src_doc_type_raw, order_no, line_no, task_no, item_code, qty_planned, qty_completed, plan_finish_date, release_time, closed_flag, closed_time, approved_flag, void_flag, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT s.tenant_id, IFNULL(s.factory_id, 1), 'AIDOP_NATIVE', 'PLATFORM_FORM', @domain, 'SALES_ORDER', IFNULL(s.source_table,'mdp_std_so'), s.order_no, s.order_line, s.order_no, s.item_code, s.order_qty, COALESCE(d.shipped_qty, s.delivered_qty), COALESCE(s.promised_delivery_date, s.plan_delivery_date, s.customer_request_date), s.order_date, IF(IFNULL(s.closed,0)=1 OR IFNULL(s.line_closed_flag,0)=1 OR d.shipped_qty>=s.order_qty, 1, 0), COALESCE(s.line_closed_time, IF(d.shipped_qty>=s.order_qty, d.last_ship_time, NULL)), 1, IFNULL(s.deleted_flag, 0), IFNULL(s.source_row_id, CAST(s.id AS CHAR)), LEFT(CONCAT('SO:', IFNULL(s.source_biz_key, s.order_no)), 200), @batch, NOW() FROM mdp_std_so s LEFT JOIN ({ShipmentActualsSql}) d ON d.tenant_id=s.tenant_id AND d.ref_task_no=s.order_no AND d.item_num=s.item_code WHERE s.tenant_id=@tid AND s.order_no IS NOT NULL AND s.order_no<>'' AND s.source_system IN ('AIDOP_NATIVE','AIDOP') ON DUPLICATE KEY UPDATE item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), qty_completed=VALUES(qty_completed), plan_finish_date=VALUES(plan_finish_date), release_time=VALUES(release_time), closed_flag=VALUES(closed_flag), closed_time=VALUES(closed_time), void_flag=VALUES(void_flag), written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time) """, new { tid = tenantId, batch = batchId, domain }); } /// /// 用销售发货实绩刷新已有 SALES_ORDER 中立行的已交货量与关行状态。 /// /// 发货流水在 S5/S7 重算时才产出,而销售订单行投影只在 S1/S2 重算时跑; /// 若不在产出发货流水的同一次跑批里回写,S7 订单发货周期要等下一次 S1 重算才有值。 /// 故本方法由 InventoryMdpSyncService 在发货投影之后调用,只更新既有行、不新增行。 /// /// public async Task RefreshSalesLineCompletionAsync( long tenantId, string? batchId, CancellationToken cancellationToken = default) { if (tenantId <= 0) return; if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_SALES", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken)) return; cancellationToken.ThrowIfCancellationRequested(); await _db.Ado.ExecuteCommandAsync( $""" UPDATE mdp_std_work_order_line l INNER JOIN ({ShipmentActualsSql}) d ON d.tenant_id=l.tenant_id AND d.ref_task_no=l.task_no AND d.item_num=l.item_code SET l.qty_completed=d.shipped_qty, l.closed_flag=IF(l.closed_flag=1 OR d.shipped_qty>=l.qty_planned, 1, 0), l.closed_time=COALESCE(l.closed_time, IF(d.shipped_qty>=l.qty_planned, d.last_ship_time, NULL)), l.sync_time=NOW() WHERE l.tenant_id=@tid AND l.doc_type='SALES_ORDER' AND l.source_system='AIDOP_NATIVE' """, new { tid = tenantId }); } }