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;
}
}