using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform; /// /// 第三方推送贴源行进入中立层。只投影已登记来源,未登记的实体直接跳过。 /// public sealed class InboundNeutralProjectionService : ITransient { private readonly ISqlSugarClient _db; private readonly MdpNeutralSourceGate _gate; public InboundNeutralProjectionService(ISqlSugarClient db, MdpNeutralSourceGate gate) { _db = db; _gate = gate; } public async Task ProjectAsync(string entityCode, long tenantId, string sourceCode, CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); var code = (entityCode ?? "").Trim().ToUpperInvariant(); var neutral = MdpSourceIdentity.NeutralCode(null, sourceCode); switch (code) { case "S6_WORK_ORDER_LINE": if (await _gate.AllowsAsync(tenantId, "WO_LINE_PROD", neutral, cancellationToken: cancellationToken)) await ProjectWorkOrderAsync(tenantId, neutral, "PROD_TASK", "mdp_stg_wo_line", cancellationToken); break; case "S7_SALES_ORDER_LINE": if (await _gate.AllowsAsync(tenantId, "WO_LINE_SALES", neutral, cancellationToken: cancellationToken)) await ProjectWorkOrderAsync(tenantId, neutral, "SALES_ORDER", "mdp_stg_wo_line", cancellationToken); break; case "S6_REPORT_TXN": if (await _gate.AllowsAsync(tenantId, "S6_REPORT", neutral, cancellationToken: cancellationToken)) await ProjectReportAsync(tenantId, neutral, cancellationToken); break; case "S5_WORK_ORDER_BOM": if (await _gate.AllowsAsync(tenantId, "WO_BOM", neutral, cancellationToken: cancellationToken)) await ProjectBomAsync(tenantId, neutral, cancellationToken); break; case "S7_FQC_TASK_TXN": if (await _gate.AllowsAsync(tenantId, "FQC_TASK", neutral, cancellationToken: cancellationToken)) await ProjectFqcAsync(tenantId, neutral, cancellationToken); break; case "S5_INVENTORY_BALANCE_MONTHLY": if (await _gate.AllowsAsync(tenantId, "INV_BAL_MONTHLY", neutral, cancellationToken: cancellationToken)) await ProjectBalanceMonthlyAsync(tenantId, neutral, cancellationToken); break; case "S2_WORK_ORDER_SCHEDULE": if (await _gate.AllowsAsync(tenantId, "WO_SCHEDULE", neutral, cancellationToken: cancellationToken)) await ProjectScheduleAsync(tenantId, neutral, cancellationToken); break; case "S3_PURCHASE_ORDER": if (await _gate.AllowsAsync(tenantId, "PURCHASE", neutral, cancellationToken: cancellationToken)) await ProjectPurchaseOrderAsync(tenantId, neutral, cancellationToken); break; case "S3_PURCHASE_RECEIPT": if (await _gate.AllowsAsync(tenantId, "PURCHASE", neutral, cancellationToken: cancellationToken)) await ProjectPurchaseReceiptAsync(tenantId, neutral, cancellationToken); break; case "S4_SHIPMENT": case "S4_IQC": case "S4_RETURN": case "S4_SHORTAGE": if (await _gate.AllowsAsync(tenantId, "SUPPLIER_DELIVERY", neutral, cancellationToken: cancellationToken)) await ProjectSupplierDeliveryAsync(code, tenantId, neutral, cancellationToken); break; case "S5_INVENTORY_TXN": if (await _gate.AllowsAsync(tenantId, "INV_TRANS", neutral, cancellationToken: cancellationToken)) await _db.Ado.ExecuteCommandAsync(ProjectApiInventorySql(), new SugarParameter("@tid", tenantId), new SugarParameter("@src", neutral)); break; } } private async Task ProjectWorkOrderAsync(long tenantId, string source, string docType, string staging, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync(staging)) return; 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, plan_finish_date, closed_flag, void_flag, approved_flag, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, 'API_INBOUND', IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''), @doc, @doc, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderNo')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LineNo')),'null') AS SIGNED), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyPlanned')),'null') AS DECIMAL(18,6)), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PlanFinishDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ClosedFlag')),'0') IN ('1','true'), 1, 0), IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.VoidFlag')),'0') IN ('1','true'), 1, 0), 1, IFNULL(s.source_row_id, s.source_biz_key), s.source_biz_key, IFNULL(s.sync_batch_id,''), NOW() FROM {staging} s WHERE s.tenant_id=@tid AND IFNULL(s.process_status,'PENDING')<>'DONE' ON DUPLICATE KEY UPDATE item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), sync_time=VALUES(sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source), new SugarParameter("@doc", docType)); } private async Task ProjectReportAsync(long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync("mdp_stg_s6_report")) return; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_s6_report (tenant_id, factory_id, source_system, written_by, work_order_no, report_date, report_qty, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, 'API_INBOUND', NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrderNo')),'null'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReportDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReportQty')),'null') AS DECIMAL(18,6)), IFNULL(s.source_row_id, s.source_biz_key), s.source_biz_key, IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_s6_report s WHERE s.tenant_id=@tid AND s.source_table='S6_REPORT_TXN' ON DUPLICATE KEY UPDATE report_qty=VALUES(report_qty), sync_time=VALUES(sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private async Task ProjectScheduleAsync(long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync("mdp_stg_schedule")) return; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_work_order_schedule (tenant_id, factory_id, source_system, doc_type, work_order, item_code, qty_ordered, qty_completed, due_date, status, prod_line, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DocType')),'null'),'PROD_TASK'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrder')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyOrdered')),'null') AS DECIMAL(18,6)), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyCompleted')),'null') AS DECIMAL(18,6)), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DueDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Status')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdLine')),'null'), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_schedule s WHERE s.tenant_id=@tid AND s.source_table='S2_WORK_ORDER_SCHEDULE' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrder')),'') NOT IN ('','null') ON DUPLICATE KEY UPDATE item_code=IF(source_system=VALUES(source_system), VALUES(item_code), item_code), qty_ordered=IF(source_system=VALUES(source_system), VALUES(qty_ordered), qty_ordered), qty_completed=IF(source_system=VALUES(source_system), VALUES(qty_completed), qty_completed), due_date=IF(source_system=VALUES(source_system), VALUES(due_date), due_date), sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private async Task ProjectPurchaseOrderAsync(long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync("mdp_stg_purchase_order")) return; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_purchase_order (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code, order_qty, due_date, order_date, status, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderQty')),'null') AS DECIMAL(18,6)), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DueDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Status')),'null'), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_purchase_order s WHERE s.tenant_id=@tid AND s.source_table='S3_PURCHASE_ORDER' ON DUPLICATE KEY UPDATE order_qty=IF(source_system=VALUES(source_system), VALUES(order_qty), order_qty), status=IF(source_system=VALUES(source_system), VALUES(status), status), sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private async Task ProjectPurchaseReceiptAsync(long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync("mdp_stg_purchase_receipt")) return; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_purchase_receipt (tenant_id, factory_id, source_system, domain, receiver, line, item_num, qty_received, pur_ord, supp, rct_date, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Receiver')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Line')),'null') AS SIGNED), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyReceived')),'null') AS DECIMAL(18,6)), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReceiptDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_purchase_receipt s WHERE s.tenant_id=@tid AND s.source_table='S3_PURCHASE_RECEIPT' ON DUPLICATE KEY UPDATE qty_received=IF(source_system=VALUES(source_system), VALUES(qty_received), qty_received), sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private async Task ProjectSupplierDeliveryAsync(string entityCode, long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); var (table, std, sql) = entityCode switch { "S4_IQC" => ("mdp_stg_s4_iqc", "mdp_std_s4_iqc", """ INSERT INTO mdp_std_s4_iqc (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code, receipt_qty, defect_qty, qc_result, receipt_date, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReceiptQty')),'null') AS DECIMAL(18,6)), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DefectQty')),'null') AS DECIMAL(18,6)), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QcResult')),'null'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReceiptDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_s4_iqc s WHERE s.tenant_id=@tid AND s.source_table='S4_IQC' ON DUPLICATE KEY UPDATE receipt_qty=IF(source_system=VALUES(source_system), VALUES(receipt_qty), receipt_qty), defect_qty=IF(source_system=VALUES(source_system), VALUES(defect_qty), defect_qty), sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time) """), "S4_RETURN" => ("mdp_stg_s4_return", "mdp_std_s4_return", """ INSERT INTO mdp_std_s4_return (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code, return_qty, return_reason, return_status, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReturnQty')),'null') AS DECIMAL(18,6)), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReturnReason')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReturnStatus')),'null'), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_s4_return s WHERE s.tenant_id=@tid AND s.source_table='S4_RETURN' ON DUPLICATE KEY UPDATE return_qty=IF(source_system=VALUES(source_system), VALUES(return_qty), return_qty), sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time) """), "S4_SHORTAGE" => ("mdp_stg_s4_shortage", "mdp_std_s4_shortage", """ INSERT INTO mdp_std_s4_shortage (tenant_id, factory_id, source_system, work_order, supplier_code, item_code, shortage_qty, risk_level, need_date, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrder')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShortageQty')),'null') AS DECIMAL(18,6)), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.RiskLevel')),'null'),'MEDIUM'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.NeedDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_s4_shortage s WHERE s.tenant_id=@tid AND s.source_table='S4_SHORTAGE' ON DUPLICATE KEY UPDATE shortage_qty=IF(source_system=VALUES(source_system), VALUES(shortage_qty), shortage_qty), sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time) """), _ => ("mdp_stg_s4_shipment", "mdp_std_s4_shipment", """ INSERT INTO mdp_std_s4_shipment (tenant_id, factory_id, source_system, shipment_no, po_no, po_line, supplier_code, item_code, ship_qty, ship_date, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipmentNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipQty')),'null') AS DECIMAL(18,6)), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_s4_shipment s WHERE s.tenant_id=@tid AND s.source_table='S4_SHIPMENT' ON DUPLICATE KEY UPDATE ship_qty=IF(source_system=VALUES(source_system), VALUES(ship_qty), ship_qty), ship_date=IF(source_system=VALUES(source_system), VALUES(ship_date), ship_date), sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time) """) }; _ = std; if (!await TableExistsAsync(table)) return; await _db.Ado.ExecuteCommandAsync(sql, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private async Task ProjectBomAsync(long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync("mdp_stg_wo_bom")) return; 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, item_code, qty_required, unit, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, 'API_INBOUND', IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LineNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyRequired')),'null') AS DECIMAL(18,6)), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Unit')),'null'), IFNULL(s.source_row_id, s.source_biz_key), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_wo_bom s WHERE s.tenant_id=@tid AND s.source_table='S5_WORK_ORDER_BOM' AND IFNULL(s.process_status,'PENDING')<>'DONE' ON DUPLICATE KEY UPDATE item_code=VALUES(item_code), qty_required=VALUES(qty_required), unit=VALUES(unit), sync_time=VALUES(sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private async Task ProjectFqcAsync(long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync("mdp_stg_fqc_pull")) return; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_fqc_task (tenant_id, source_system, written_by, bill_no, production_order_no, material_code, qty, sales_order_no, domain, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT @tid, @src, 'API_INBOUND', NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BillNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProductionOrderNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.MaterialCode')),'null'), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Qty')),'null') AS DECIMAL(18,6)), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SalesOrderNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'), IFNULL(s.source_row_id, s.source_biz_key), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_fqc_pull s WHERE s.tenant_id=@tid AND s.source_table='S7_FQC_TASK_TXN' AND IFNULL(s.process_status,'PENDING')<>'DONE' ON DUPLICATE KEY UPDATE material_code=VALUES(material_code), qty=VALUES(qty), sync_time=VALUES(sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private async Task ProjectBalanceMonthlyAsync(long tenantId, string source, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); if (!await TableExistsAsync("mdp_stg_inv_bal_monthly")) return; await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_inventory_balance_monthly (tenant_id, factory_id, source_system, written_by, period_ym, category_code, category_name, warehouse_code, warehouse_name, item_code, avg_balance_amount, issue_cost_amount, source_biz_key, sync_batch_id, sync_time) SELECT @tid, 1, @src, 'API_INBOUND', NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PeriodYm')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CategoryCode')),'null'),''), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CategoryName')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WarehouseCode')),'null'),''), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WarehouseName')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.AvgBalanceAmount')),'null') AS DECIMAL(18,6)), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.IssueCostAmount')),'null') AS DECIMAL(18,6)), IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW() FROM mdp_stg_inv_bal_monthly s WHERE s.tenant_id=@tid AND s.source_table='S5_INVENTORY_BALANCE_MONTHLY' AND IFNULL(s.process_status,'PENDING')<>'DONE' ON DUPLICATE KEY UPDATE avg_balance_amount=VALUES(avg_balance_amount), issue_cost_amount=VALUES(issue_cost_amount), sync_time=VALUES(sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source)); } private static string ProjectApiInventorySql() => """ UPDATE mdp_std_inv_trans SET written_by='API_INBOUND' WHERE tenant_id=@tid AND source_system=@src AND written_by IS NULL """; private async Task TableExistsAsync(string table) { var n = await _db.Ado.GetIntAsync( "SELECT COUNT(*) FROM information_schema.TABLES WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t", new SugarParameter("@t", table)); return n > 0; } }