InboundNeutralProjectionService.cs 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422
  1. using SqlSugar;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  3. /// <summary>
  4. /// 第三方推送贴源行进入中立层。只投影已登记来源,未登记的实体直接跳过。
  5. /// </summary>
  6. public sealed class InboundNeutralProjectionService : ITransient
  7. {
  8. private readonly ISqlSugarClient _db;
  9. private readonly MdpNeutralSourceGate _gate;
  10. public InboundNeutralProjectionService(ISqlSugarClient db, MdpNeutralSourceGate gate)
  11. {
  12. _db = db;
  13. _gate = gate;
  14. }
  15. public async Task ProjectAsync(string entityCode, long tenantId, string sourceCode, CancellationToken cancellationToken = default)
  16. {
  17. cancellationToken.ThrowIfCancellationRequested();
  18. var code = (entityCode ?? "").Trim().ToUpperInvariant();
  19. var neutral = MdpSourceIdentity.NeutralCode(null, sourceCode);
  20. switch (code)
  21. {
  22. case "S6_WORK_ORDER_LINE":
  23. if (await _gate.AllowsAsync(tenantId, "WO_LINE_PROD", neutral, cancellationToken: cancellationToken))
  24. await ProjectWorkOrderAsync(tenantId, neutral, "PROD_TASK", "mdp_stg_wo_line", cancellationToken);
  25. break;
  26. case "S7_SALES_ORDER_LINE":
  27. if (await _gate.AllowsAsync(tenantId, "WO_LINE_SALES", neutral, cancellationToken: cancellationToken))
  28. await ProjectWorkOrderAsync(tenantId, neutral, "SALES_ORDER", "mdp_stg_wo_line", cancellationToken);
  29. break;
  30. case "S6_REPORT_TXN":
  31. if (await _gate.AllowsAsync(tenantId, "S6_REPORT", neutral, cancellationToken: cancellationToken))
  32. await ProjectReportAsync(tenantId, neutral, cancellationToken);
  33. break;
  34. case "S5_WORK_ORDER_BOM":
  35. if (await _gate.AllowsAsync(tenantId, "WO_BOM", neutral, cancellationToken: cancellationToken))
  36. await ProjectBomAsync(tenantId, neutral, cancellationToken);
  37. break;
  38. case "S7_FQC_TASK_TXN":
  39. if (await _gate.AllowsAsync(tenantId, "FQC_TASK", neutral, cancellationToken: cancellationToken))
  40. await ProjectFqcAsync(tenantId, neutral, cancellationToken);
  41. break;
  42. case "S5_INVENTORY_BALANCE_MONTHLY":
  43. if (await _gate.AllowsAsync(tenantId, "INV_BAL_MONTHLY", neutral, cancellationToken: cancellationToken))
  44. await ProjectBalanceMonthlyAsync(tenantId, neutral, cancellationToken);
  45. break;
  46. case "S2_WORK_ORDER_SCHEDULE":
  47. if (await _gate.AllowsAsync(tenantId, "WO_SCHEDULE", neutral, cancellationToken: cancellationToken))
  48. await ProjectScheduleAsync(tenantId, neutral, cancellationToken);
  49. break;
  50. case "S3_PURCHASE_ORDER":
  51. if (await _gate.AllowsAsync(tenantId, "PURCHASE", neutral, cancellationToken: cancellationToken))
  52. await ProjectPurchaseOrderAsync(tenantId, neutral, cancellationToken);
  53. break;
  54. case "S3_PURCHASE_RECEIPT":
  55. if (await _gate.AllowsAsync(tenantId, "PURCHASE", neutral, cancellationToken: cancellationToken))
  56. await ProjectPurchaseReceiptAsync(tenantId, neutral, cancellationToken);
  57. break;
  58. case "S4_SHIPMENT":
  59. case "S4_IQC":
  60. case "S4_RETURN":
  61. case "S4_SHORTAGE":
  62. if (await _gate.AllowsAsync(tenantId, "SUPPLIER_DELIVERY", neutral, cancellationToken: cancellationToken))
  63. await ProjectSupplierDeliveryAsync(code, tenantId, neutral, cancellationToken);
  64. break;
  65. case "S5_INVENTORY_TXN":
  66. if (await _gate.AllowsAsync(tenantId, "INV_TRANS", neutral, cancellationToken: cancellationToken))
  67. await _db.Ado.ExecuteCommandAsync(ProjectApiInventorySql(), new SugarParameter("@tid", tenantId), new SugarParameter("@src", neutral));
  68. break;
  69. }
  70. }
  71. private async Task ProjectWorkOrderAsync(long tenantId, string source, string docType, string staging, CancellationToken cancellationToken)
  72. {
  73. cancellationToken.ThrowIfCancellationRequested();
  74. if (!await TableExistsAsync(staging)) return;
  75. await _db.Ado.ExecuteCommandAsync(
  76. $"""
  77. INSERT INTO mdp_std_work_order_line
  78. (tenant_id, factory_id, source_system, written_by, domain, doc_type, src_doc_type_raw,
  79. order_no, line_no, task_no, item_code, qty_planned, plan_finish_date,
  80. closed_flag, void_flag, approved_flag, source_row_id, source_biz_key, sync_batch_id, sync_time)
  81. SELECT @tid, 1, @src, 'API_INBOUND',
  82. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''),
  83. @doc, @doc,
  84. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderNo')),'null'),
  85. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LineNo')),'null') AS SIGNED),
  86. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderNo')),'null'),
  87. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),
  88. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyPlanned')),'null') AS DECIMAL(18,6)),
  89. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PlanFinishDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  90. IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ClosedFlag')),'0') IN ('1','true'), 1, 0),
  91. IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.VoidFlag')),'0') IN ('1','true'), 1, 0),
  92. 1,
  93. IFNULL(s.source_row_id, s.source_biz_key), s.source_biz_key, IFNULL(s.sync_batch_id,''), NOW()
  94. FROM {staging} s
  95. WHERE s.tenant_id=@tid AND IFNULL(s.process_status,'PENDING')<>'DONE'
  96. ON DUPLICATE KEY UPDATE
  97. item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), sync_time=VALUES(sync_time)
  98. """,
  99. new SugarParameter("@tid", tenantId),
  100. new SugarParameter("@src", source),
  101. new SugarParameter("@doc", docType));
  102. }
  103. private async Task ProjectReportAsync(long tenantId, string source, CancellationToken cancellationToken)
  104. {
  105. cancellationToken.ThrowIfCancellationRequested();
  106. if (!await TableExistsAsync("mdp_stg_s6_report")) return;
  107. await _db.Ado.ExecuteCommandAsync(
  108. """
  109. INSERT INTO mdp_std_s6_report
  110. (tenant_id, factory_id, source_system, written_by, work_order_no, report_date, report_qty,
  111. source_row_id, source_biz_key, sync_batch_id, sync_time)
  112. SELECT @tid, 1, @src, 'API_INBOUND',
  113. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrderNo')),'null'),
  114. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReportDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  115. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReportQty')),'null') AS DECIMAL(18,6)),
  116. IFNULL(s.source_row_id, s.source_biz_key), s.source_biz_key, IFNULL(s.sync_batch_id,''), NOW()
  117. FROM mdp_stg_s6_report s
  118. WHERE s.tenant_id=@tid AND s.source_table='S6_REPORT_TXN'
  119. ON DUPLICATE KEY UPDATE report_qty=VALUES(report_qty), sync_time=VALUES(sync_time)
  120. """,
  121. new SugarParameter("@tid", tenantId),
  122. new SugarParameter("@src", source));
  123. }
  124. private async Task ProjectScheduleAsync(long tenantId, string source, CancellationToken cancellationToken)
  125. {
  126. cancellationToken.ThrowIfCancellationRequested();
  127. if (!await TableExistsAsync("mdp_stg_schedule")) return;
  128. await _db.Ado.ExecuteCommandAsync(
  129. """
  130. INSERT INTO mdp_std_work_order_schedule
  131. (tenant_id, factory_id, source_system, doc_type, work_order, item_code,
  132. qty_ordered, qty_completed, due_date, status, prod_line,
  133. source_biz_key, sync_batch_id, sync_time)
  134. SELECT @tid, 1, @src,
  135. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DocType')),'null'),'PROD_TASK'),
  136. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrder')),'null'),
  137. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),
  138. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyOrdered')),'null') AS DECIMAL(18,6)),
  139. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyCompleted')),'null') AS DECIMAL(18,6)),
  140. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DueDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  141. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Status')),'null'),
  142. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdLine')),'null'),
  143. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')),
  144. IFNULL(s.sync_batch_id,''), NOW()
  145. FROM mdp_stg_schedule s
  146. WHERE s.tenant_id=@tid AND s.source_table='S2_WORK_ORDER_SCHEDULE'
  147. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrder')),'') NOT IN ('','null')
  148. ON DUPLICATE KEY UPDATE
  149. item_code=IF(source_system=VALUES(source_system), VALUES(item_code), item_code),
  150. qty_ordered=IF(source_system=VALUES(source_system), VALUES(qty_ordered), qty_ordered),
  151. qty_completed=IF(source_system=VALUES(source_system), VALUES(qty_completed), qty_completed),
  152. due_date=IF(source_system=VALUES(source_system), VALUES(due_date), due_date),
  153. sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time)
  154. """,
  155. new SugarParameter("@tid", tenantId),
  156. new SugarParameter("@src", source));
  157. }
  158. private async Task ProjectPurchaseOrderAsync(long tenantId, string source, CancellationToken cancellationToken)
  159. {
  160. cancellationToken.ThrowIfCancellationRequested();
  161. if (!await TableExistsAsync("mdp_stg_purchase_order")) return;
  162. await _db.Ado.ExecuteCommandAsync(
  163. """
  164. INSERT INTO mdp_std_purchase_order
  165. (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code,
  166. order_qty, due_date, order_date, status, source_biz_key, sync_batch_id, sync_time)
  167. SELECT @tid, 1, @src,
  168. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'),
  169. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'),
  170. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),
  171. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),
  172. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderQty')),'null') AS DECIMAL(18,6)),
  173. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DueDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  174. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  175. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Status')),'null'),
  176. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')),
  177. IFNULL(s.sync_batch_id,''), NOW()
  178. FROM mdp_stg_purchase_order s
  179. WHERE s.tenant_id=@tid AND s.source_table='S3_PURCHASE_ORDER'
  180. ON DUPLICATE KEY UPDATE
  181. order_qty=IF(source_system=VALUES(source_system), VALUES(order_qty), order_qty),
  182. status=IF(source_system=VALUES(source_system), VALUES(status), status),
  183. sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time)
  184. """,
  185. new SugarParameter("@tid", tenantId),
  186. new SugarParameter("@src", source));
  187. }
  188. private async Task ProjectPurchaseReceiptAsync(long tenantId, string source, CancellationToken cancellationToken)
  189. {
  190. cancellationToken.ThrowIfCancellationRequested();
  191. if (!await TableExistsAsync("mdp_stg_purchase_receipt")) return;
  192. await _db.Ado.ExecuteCommandAsync(
  193. """
  194. INSERT INTO mdp_std_purchase_receipt
  195. (tenant_id, factory_id, source_system, domain, receiver, line, item_num,
  196. qty_received, pur_ord, supp, rct_date, source_biz_key, sync_batch_id, sync_time)
  197. SELECT @tid, 1, @src,
  198. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''),
  199. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Receiver')),'null'),
  200. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Line')),'null') AS SIGNED),
  201. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),
  202. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyReceived')),'null') AS DECIMAL(18,6)),
  203. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'),
  204. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),
  205. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReceiptDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  206. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')),
  207. IFNULL(s.sync_batch_id,''), NOW()
  208. FROM mdp_stg_purchase_receipt s
  209. WHERE s.tenant_id=@tid AND s.source_table='S3_PURCHASE_RECEIPT'
  210. ON DUPLICATE KEY UPDATE
  211. qty_received=IF(source_system=VALUES(source_system), VALUES(qty_received), qty_received),
  212. sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time)
  213. """,
  214. new SugarParameter("@tid", tenantId),
  215. new SugarParameter("@src", source));
  216. }
  217. private async Task ProjectSupplierDeliveryAsync(string entityCode, long tenantId, string source, CancellationToken cancellationToken)
  218. {
  219. cancellationToken.ThrowIfCancellationRequested();
  220. var (table, std, sql) = entityCode switch
  221. {
  222. "S4_IQC" => ("mdp_stg_s4_iqc", "mdp_std_s4_iqc", """
  223. INSERT INTO mdp_std_s4_iqc
  224. (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code,
  225. receipt_qty, defect_qty, qc_result, receipt_date, source_biz_key, sync_batch_id, sync_time)
  226. SELECT @tid, 1, @src,
  227. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'),
  228. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'),
  229. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''),
  230. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''),
  231. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReceiptQty')),'null') AS DECIMAL(18,6)),
  232. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DefectQty')),'null') AS DECIMAL(18,6)),
  233. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QcResult')),'null'),
  234. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReceiptDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  235. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW()
  236. FROM mdp_stg_s4_iqc s
  237. WHERE s.tenant_id=@tid AND s.source_table='S4_IQC'
  238. ON DUPLICATE KEY UPDATE
  239. receipt_qty=IF(source_system=VALUES(source_system), VALUES(receipt_qty), receipt_qty),
  240. defect_qty=IF(source_system=VALUES(source_system), VALUES(defect_qty), defect_qty),
  241. sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time)
  242. """),
  243. "S4_RETURN" => ("mdp_stg_s4_return", "mdp_std_s4_return", """
  244. INSERT INTO mdp_std_s4_return
  245. (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code,
  246. return_qty, return_reason, return_status, source_biz_key, sync_batch_id, sync_time)
  247. SELECT @tid, 1, @src,
  248. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'),
  249. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'),
  250. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''),
  251. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''),
  252. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReturnQty')),'null') AS DECIMAL(18,6)),
  253. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReturnReason')),'null'),
  254. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReturnStatus')),'null'),
  255. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW()
  256. FROM mdp_stg_s4_return s
  257. WHERE s.tenant_id=@tid AND s.source_table='S4_RETURN'
  258. ON DUPLICATE KEY UPDATE
  259. return_qty=IF(source_system=VALUES(source_system), VALUES(return_qty), return_qty),
  260. sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time)
  261. """),
  262. "S4_SHORTAGE" => ("mdp_stg_s4_shortage", "mdp_std_s4_shortage", """
  263. INSERT INTO mdp_std_s4_shortage
  264. (tenant_id, factory_id, source_system, work_order, supplier_code, item_code,
  265. shortage_qty, risk_level, need_date, source_biz_key, sync_batch_id, sync_time)
  266. SELECT @tid, 1, @src,
  267. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrder')),'null'),
  268. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''),
  269. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''),
  270. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShortageQty')),'null') AS DECIMAL(18,6)),
  271. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.RiskLevel')),'null'),'MEDIUM'),
  272. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.NeedDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  273. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW()
  274. FROM mdp_stg_s4_shortage s
  275. WHERE s.tenant_id=@tid AND s.source_table='S4_SHORTAGE'
  276. ON DUPLICATE KEY UPDATE
  277. shortage_qty=IF(source_system=VALUES(source_system), VALUES(shortage_qty), shortage_qty),
  278. sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time)
  279. """),
  280. _ => ("mdp_stg_s4_shipment", "mdp_std_s4_shipment", """
  281. INSERT INTO mdp_std_s4_shipment
  282. (tenant_id, factory_id, source_system, shipment_no, po_no, po_line, supplier_code, item_code,
  283. ship_qty, ship_date, source_biz_key, sync_batch_id, sync_time)
  284. SELECT @tid, 1, @src,
  285. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipmentNo')),'null'),
  286. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoNo')),'null'),
  287. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PoLine')),'null'),
  288. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SupplierCode')),'null'),''),
  289. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''),
  290. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipQty')),'null') AS DECIMAL(18,6)),
  291. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'),
  292. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')), IFNULL(s.sync_batch_id,''), NOW()
  293. FROM mdp_stg_s4_shipment s
  294. WHERE s.tenant_id=@tid AND s.source_table='S4_SHIPMENT'
  295. ON DUPLICATE KEY UPDATE
  296. ship_qty=IF(source_system=VALUES(source_system), VALUES(ship_qty), ship_qty),
  297. ship_date=IF(source_system=VALUES(source_system), VALUES(ship_date), ship_date),
  298. sync_time=IF(source_system=VALUES(source_system), VALUES(sync_time), sync_time)
  299. """)
  300. };
  301. _ = std;
  302. if (!await TableExistsAsync(table)) return;
  303. await _db.Ado.ExecuteCommandAsync(sql, new SugarParameter("@tid", tenantId), new SugarParameter("@src", source));
  304. }
  305. private async Task ProjectBomAsync(long tenantId, string source, CancellationToken cancellationToken)
  306. {
  307. cancellationToken.ThrowIfCancellationRequested();
  308. if (!await TableExistsAsync("mdp_stg_wo_bom")) return;
  309. await _db.Ado.ExecuteCommandAsync(
  310. """
  311. INSERT INTO mdp_std_work_order_bom
  312. (tenant_id, factory_id, source_system, written_by, domain, order_no, order_line_no,
  313. item_code, qty_required, unit, source_row_id, source_biz_key, sync_batch_id, sync_time)
  314. SELECT @tid, 1, @src, 'API_INBOUND',
  315. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''),
  316. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrderNo')),'null'),
  317. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LineNo')),'null'),
  318. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),
  319. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyRequired')),'null') AS DECIMAL(18,6)),
  320. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Unit')),'null'),
  321. IFNULL(s.source_row_id, s.source_biz_key),
  322. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')),
  323. IFNULL(s.sync_batch_id,''), NOW()
  324. FROM mdp_stg_wo_bom s
  325. WHERE s.tenant_id=@tid AND s.source_table='S5_WORK_ORDER_BOM'
  326. AND IFNULL(s.process_status,'PENDING')<>'DONE'
  327. ON DUPLICATE KEY UPDATE
  328. item_code=VALUES(item_code), qty_required=VALUES(qty_required),
  329. unit=VALUES(unit), sync_time=VALUES(sync_time)
  330. """,
  331. new SugarParameter("@tid", tenantId),
  332. new SugarParameter("@src", source));
  333. }
  334. private async Task ProjectFqcAsync(long tenantId, string source, CancellationToken cancellationToken)
  335. {
  336. cancellationToken.ThrowIfCancellationRequested();
  337. if (!await TableExistsAsync("mdp_stg_fqc_pull")) return;
  338. await _db.Ado.ExecuteCommandAsync(
  339. """
  340. INSERT INTO mdp_std_fqc_task
  341. (tenant_id, source_system, written_by, bill_no, production_order_no, material_code, qty,
  342. sales_order_no, domain, source_row_id, source_biz_key, sync_batch_id, sync_time)
  343. SELECT @tid, @src, 'API_INBOUND',
  344. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BillNo')),'null'),
  345. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProductionOrderNo')),'null'),
  346. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.MaterialCode')),'null'),
  347. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Qty')),'null') AS DECIMAL(18,6)),
  348. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SalesOrderNo')),'null'),
  349. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),
  350. IFNULL(s.source_row_id, s.source_biz_key),
  351. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')),
  352. IFNULL(s.sync_batch_id,''), NOW()
  353. FROM mdp_stg_fqc_pull s
  354. WHERE s.tenant_id=@tid AND s.source_table='S7_FQC_TASK_TXN'
  355. AND IFNULL(s.process_status,'PENDING')<>'DONE'
  356. ON DUPLICATE KEY UPDATE
  357. material_code=VALUES(material_code), qty=VALUES(qty), sync_time=VALUES(sync_time)
  358. """,
  359. new SugarParameter("@tid", tenantId),
  360. new SugarParameter("@src", source));
  361. }
  362. private async Task ProjectBalanceMonthlyAsync(long tenantId, string source, CancellationToken cancellationToken)
  363. {
  364. cancellationToken.ThrowIfCancellationRequested();
  365. if (!await TableExistsAsync("mdp_stg_inv_bal_monthly")) return;
  366. await _db.Ado.ExecuteCommandAsync(
  367. """
  368. INSERT INTO mdp_std_inventory_balance_monthly
  369. (tenant_id, factory_id, source_system, written_by, period_ym,
  370. category_code, category_name, warehouse_code, warehouse_name, item_code,
  371. avg_balance_amount, issue_cost_amount, source_biz_key, sync_batch_id, sync_time)
  372. SELECT @tid, 1, @src, 'API_INBOUND',
  373. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.PeriodYm')),'null'),
  374. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CategoryCode')),'null'),''),
  375. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CategoryName')),'null'),
  376. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WarehouseCode')),'null'),''),
  377. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WarehouseName')),'null'),
  378. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'),''),
  379. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.AvgBalanceAmount')),'null') AS DECIMAL(18,6)),
  380. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.IssueCostAmount')),'null') AS DECIMAL(18,6)),
  381. IFNULL(s.source_biz_key, IFNULL(s.source_row_id,'')),
  382. IFNULL(s.sync_batch_id,''), NOW()
  383. FROM mdp_stg_inv_bal_monthly s
  384. WHERE s.tenant_id=@tid AND s.source_table='S5_INVENTORY_BALANCE_MONTHLY'
  385. AND IFNULL(s.process_status,'PENDING')<>'DONE'
  386. ON DUPLICATE KEY UPDATE
  387. avg_balance_amount=VALUES(avg_balance_amount),
  388. issue_cost_amount=VALUES(issue_cost_amount),
  389. sync_time=VALUES(sync_time)
  390. """,
  391. new SugarParameter("@tid", tenantId),
  392. new SugarParameter("@src", source));
  393. }
  394. private static string ProjectApiInventorySql() =>
  395. """
  396. UPDATE mdp_std_inv_trans
  397. SET written_by='API_INBOUND'
  398. WHERE tenant_id=@tid AND source_system=@src AND written_by IS NULL
  399. """;
  400. private async Task<bool> TableExistsAsync(string table)
  401. {
  402. var n = await _db.Ado.GetIntAsync(
  403. "SELECT COUNT(*) FROM information_schema.TABLES WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t",
  404. new SugarParameter("@t", table));
  405. return n > 0;
  406. }
  407. }