NativeNeutralProjectionService.cs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Schema;
  2. using Admin.NET.Plugin.AiDOP.Infrastructure;
  3. using Microsoft.Extensions.Logging;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  5. /// <summary>
  6. /// 自有业务表投影到中立层。只写登记为 AIDOP_NATIVE 的对象。
  7. /// </summary>
  8. public sealed class NativeNeutralProjectionService : ITransient
  9. {
  10. private readonly ISqlSugarClient _db;
  11. private readonly MdpNeutralSourceGate _neutralGate;
  12. private readonly EmployeePositionMapService _positions;
  13. private readonly SourceDomainTenantResolver _domains;
  14. private readonly ILogger _logger;
  15. public NativeNeutralProjectionService(
  16. ISqlSugarClient db,
  17. MdpNeutralSourceGate neutralGate,
  18. EmployeePositionMapService positions,
  19. SourceDomainTenantResolver domains,
  20. ILoggerFactory loggerFactory)
  21. {
  22. _db = db;
  23. _neutralGate = neutralGate;
  24. _positions = positions;
  25. _domains = domains;
  26. _logger = loggerFactory.CreateLogger(nameof(NativeNeutralProjectionService));
  27. }
  28. /// <summary>
  29. /// 销售发货实绩(登记来源的 <c>SALES_SHIP</c>)按「订单号 + 物料」汇总。
  30. /// <para>
  31. /// 模式一下销售订单的权威在自有 S1,执行事实由外部系统提供,故订单行的已交货量与关行状态
  32. /// 要用发货流水推导,不能只看 <c>mdp_std_so.delivered_qty</c>(自建单常年为空)。
  33. /// 来源不写死:内联 <c>mdp_tenant_std_source</c>,只读该租户登记为 INV_TRANS 权威的那一个来源。
  34. /// 粒度与 S7 指标 SQL 的 JOIN 一致(中立流水没有订单行号,最细只能到订单号 + 物料)。
  35. /// </para>
  36. /// </summary>
  37. private const string ShipmentActualsSql = """
  38. SELECT t.tenant_id, t.ref_task_no, t.item_num,
  39. SUM(ABS(IFNULL(t.qty_change,0))) AS shipped_qty,
  40. MAX(t.approved_time) AS last_ship_time
  41. FROM mdp_std_inv_trans t
  42. INNER JOIN mdp_tenant_std_source ts
  43. ON ts.tenant_id=t.tenant_id AND ts.std_object='INV_TRANS' AND ts.source_system=t.source_system
  44. WHERE t.tenant_id=@tid AND t.biz_doc_type='SALES_SHIP'
  45. AND t.summary_flag=0 AND t.void_flag=0 AND t.approved_flag=1
  46. GROUP BY t.tenant_id, t.ref_task_no, t.item_num
  47. """;
  48. public async Task ProjectAsync(long tenantId, string batchId, CancellationToken cancellationToken = default)
  49. {
  50. if (tenantId <= 0) return;
  51. var failed = new List<string>();
  52. await Try(failed, "WO_LINE_PROD", () => ProjectWorkOrdersAsync(tenantId, batchId, cancellationToken));
  53. await Try(failed, "WO_SCHEDULE", () => ProjectWorkOrderScheduleAsync(tenantId, batchId, cancellationToken));
  54. await Try(failed, "WO_BOM", () => ProjectBomAsync(tenantId, batchId, cancellationToken));
  55. await Try(failed, "EMPLOYEE", () => ProjectEmployeesAsync(tenantId, batchId, cancellationToken));
  56. await Try(failed, "WO_LINE_SALES", () => ProjectSalesLinesAsync(tenantId, batchId, cancellationToken));
  57. if (failed.Count > 0)
  58. _logger.LogWarning("自有中立投影未完成 tenant={Tenant} {Failed}", tenantId, string.Join(" | ", failed));
  59. }
  60. private async Task Try(List<string> failed, string name, Func<Task> action)
  61. {
  62. try
  63. {
  64. await action();
  65. }
  66. catch (Exception ex) when (ex is not OperationCanceledException)
  67. {
  68. failed.Add($"{name}: {ex.Message}");
  69. }
  70. }
  71. private async Task ProjectWorkOrdersAsync(long tenantId, string batchId, CancellationToken cancellationToken)
  72. {
  73. if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_PROD", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
  74. return;
  75. await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_line");
  76. cancellationToken.ThrowIfCancellationRequested();
  77. await _db.Ado.ExecuteCommandAsync(
  78. """
  79. INSERT INTO mdp_std_work_order_line
  80. (tenant_id, factory_id, source_system, written_by, domain, doc_type, src_doc_type_raw,
  81. order_no, line_no, task_no, item_code, qty_planned, qty_completed, plan_finish_date, release_time,
  82. closed_flag, closed_time, approved_flag, void_flag, source_row_id, source_biz_key, sync_batch_id, sync_time)
  83. SELECT
  84. w.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(w.Domain,''), 'PROD_TASK', 'WorkOrdMaster',
  85. w.WorkOrd, NULL, w.WorkOrd, w.ItemNum, w.QtyOrded, w.QtyCompleted, w.DueDate, w.ReleaseDate,
  86. IF(UPPER(IFNULL(w.Status,''))='C', 1, 0), NULL, 1, IF(IFNULL(w.IsActive,1)=0, 1, 0),
  87. CAST(w.RecID AS CHAR), LEFT(CONCAT(IFNULL(w.Domain,''), ':', w.WorkOrd), 200), @batch, NOW()
  88. FROM WorkOrdMaster w
  89. WHERE w.tenant_id=@tid AND w.WorkOrd IS NOT NULL AND w.WorkOrd<>''
  90. ON DUPLICATE KEY UPDATE
  91. item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), qty_completed=VALUES(qty_completed),
  92. plan_finish_date=VALUES(plan_finish_date), release_time=VALUES(release_time),
  93. closed_flag=VALUES(closed_flag), void_flag=VALUES(void_flag),
  94. written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  95. """,
  96. new { tid = tenantId, batch = batchId });
  97. }
  98. /// <summary>
  99. /// 自有生产工单 → <c>mdp_std_work_order_schedule</c>(<c>doc_type='PROD_TASK'</c>)。
  100. /// <para>
  101. /// S5_L1_002 物料齐套满足率的分母以工单排程头为驱动,且 BOM 按 <c>source_system</c> 关联,
  102. /// 故排程必须与 <see cref="ProjectBomAsync"/> 写同一个 <c>AIDOP_NATIVE</c>,否则分母恒为空、
  103. /// 指标只会算出 NO_DATA。
  104. /// </para>
  105. /// <para>
  106. /// 唯一键 <c>uk_std_wo_sched_type</c> 不含 <c>source_system</c>,同一工单号若已被别的来源占用,
  107. /// 直接 upsert 会把对方的行改成自有来源,故与 T8 投影一样加 NOT EXISTS 反向占用保护。
  108. /// </para>
  109. /// </summary>
  110. private async Task ProjectWorkOrderScheduleAsync(long tenantId, string batchId, CancellationToken cancellationToken)
  111. {
  112. if (!await _neutralGate.AllowsAsync(tenantId, "WO_SCHEDULE", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
  113. return;
  114. await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_schedule");
  115. cancellationToken.ThrowIfCancellationRequested();
  116. await _db.Ado.ExecuteCommandAsync(
  117. """
  118. INSERT INTO mdp_std_work_order_schedule
  119. (tenant_id, factory_id, source_system, written_by, work_order, doc_type,
  120. item_code, item_name, site_code, status, priority, urgent_flag,
  121. qty_ordered, qty_completed, order_date, due_date, release_date,
  122. prod_line, lot_serial, drawing_no, project, work_order_type, labor_variance,
  123. approved_flag, void_flag, source_biz_key, sync_batch_id, sync_time)
  124. SELECT
  125. w.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', w.WorkOrd, 'PROD_TASK',
  126. w.ItemNum, w.ItemName, w.Site, w.Status, w.Priority, IF(IFNULL(w.Urgent,0)<>0, 1, 0),
  127. w.QtyOrded, w.QtyCompleted, w.OrdDate, w.DueDate, w.ReleaseDate,
  128. w.ProdLine, w.Batch, w.Drawing, w.Project, w.Typed, w.LbrVar,
  129. 1, IF(IFNULL(w.IsActive,1)=0, 1, 0),
  130. LEFT(CONCAT(IFNULL(w.Domain,''), ':', w.WorkOrd), 200), @batch, NOW()
  131. FROM WorkOrdMaster w
  132. WHERE w.tenant_id=@tid AND w.WorkOrd IS NOT NULL AND w.WorkOrd<>''
  133. AND NOT EXISTS (
  134. SELECT 1 FROM mdp_std_work_order_schedule x
  135. WHERE x.tenant_id=w.tenant_id AND x.doc_type='PROD_TASK' AND x.work_order=w.WorkOrd
  136. AND x.source_system<>'AIDOP_NATIVE')
  137. ON DUPLICATE KEY UPDATE
  138. item_code=VALUES(item_code), item_name=VALUES(item_name), site_code=VALUES(site_code),
  139. status=VALUES(status), priority=VALUES(priority), urgent_flag=VALUES(urgent_flag),
  140. qty_ordered=VALUES(qty_ordered), qty_completed=VALUES(qty_completed),
  141. order_date=VALUES(order_date), due_date=VALUES(due_date), release_date=VALUES(release_date),
  142. prod_line=VALUES(prod_line), lot_serial=VALUES(lot_serial), drawing_no=VALUES(drawing_no),
  143. project=VALUES(project), work_order_type=VALUES(work_order_type),
  144. labor_variance=VALUES(labor_variance),
  145. approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag),
  146. written_by=VALUES(written_by), source_biz_key=VALUES(source_biz_key),
  147. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  148. """,
  149. new { tid = tenantId, batch = batchId });
  150. }
  151. private async Task ProjectBomAsync(long tenantId, string batchId, CancellationToken cancellationToken)
  152. {
  153. if (!await _neutralGate.AllowsAsync(tenantId, "WO_BOM", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
  154. return;
  155. await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_bom");
  156. cancellationToken.ThrowIfCancellationRequested();
  157. await _db.Ado.ExecuteCommandAsync(
  158. """
  159. INSERT INTO mdp_std_work_order_bom
  160. (tenant_id, factory_id, source_system, written_by, domain, order_no, order_line_no, task_no, item_code,
  161. qty_required, unit, source_row_id, source_biz_key, sync_batch_id, sync_time)
  162. SELECT
  163. d.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(d.Domain,''), d.WorkOrd,
  164. CAST(d.LineNum AS CHAR), d.WorkOrd, d.ItemNum, d.QtyRequired, d.UM,
  165. CAST(d.RecID AS CHAR),
  166. LEFT(CONCAT(IFNULL(d.Domain,''), ':', d.WorkOrd, ':', d.LineNum, ':', IFNULL(d.ItemNum,'')), 200),
  167. @batch, NOW()
  168. FROM WorkOrdDetail d
  169. WHERE d.tenant_id=@tid AND d.WorkOrd IS NOT NULL AND d.WorkOrd<>''
  170. ON DUPLICATE KEY UPDATE
  171. item_code=VALUES(item_code), qty_required=VALUES(qty_required), unit=VALUES(unit),
  172. written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  173. """,
  174. new { tid = tenantId, batch = batchId });
  175. }
  176. private async Task ProjectEmployeesAsync(long tenantId, string batchId, CancellationToken cancellationToken)
  177. {
  178. if (!await _neutralGate.AllowsAsync(tenantId, "EMPLOYEE", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
  179. return;
  180. await _positions.SyncNativeAsync(tenantId, cancellationToken);
  181. await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_employee");
  182. cancellationToken.ThrowIfCancellationRequested();
  183. await _db.Ado.ExecuteCommandAsync(
  184. """
  185. INSERT INTO mdp_std_employee
  186. (tenant_id, factory_id, source_system, written_by, domain, employee_no, employee_name, department_code,
  187. position_code, src_position_raw, employment_status, src_employment_status_raw,
  188. source_row_id, source_biz_key, sync_batch_id, sync_time)
  189. SELECT
  190. e.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(e.Domain,''), e.Employee, e.Name, e.Department,
  191. IFNULL(m.position_code, 'UNKNOWN'),
  192. COALESCE(NULLIF(e.JobTitle,''), NULLIF(e.WorkCtr,'')),
  193. CASE
  194. WHEN e.DateTerminated IS NOT NULL AND e.DateTerminated <= CURDATE() THEN 'LEFT'
  195. WHEN e.EmploymentStatus IN ('在职') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('ACTIVE','ONJOB','ON_JOB','AC') THEN 'ACTIVE'
  196. WHEN e.EmploymentStatus IN ('离职','辞退') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('LEFT','LEAVE','RESIGNED','TERMINATED','QUIT','DC') THEN 'LEFT'
  197. WHEN e.EmploymentStatus IN ('停用') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('INACTIVE','DISABLED','SUSPENDED') THEN 'INACTIVE'
  198. ELSE 'UNKNOWN' END,
  199. e.EmploymentStatus,
  200. CAST(e.RecID AS CHAR), LEFT(CONCAT(IFNULL(e.Domain,''), ':', e.Employee), 200), @batch, NOW()
  201. FROM EmployeeMaster e
  202. LEFT JOIN mdp_employee_position_map m
  203. ON m.tenant_id=e.tenant_id AND m.source_system='AIDOP_NATIVE'
  204. AND m.domain=IFNULL(e.Domain,'')
  205. AND m.src_position_raw=COALESCE(NULLIF(e.JobTitle,''), NULLIF(e.WorkCtr,''))
  206. WHERE e.tenant_id=@tid AND e.Employee IS NOT NULL AND e.Employee<>''
  207. ON DUPLICATE KEY UPDATE
  208. employee_name=VALUES(employee_name), department_code=VALUES(department_code),
  209. position_code=VALUES(position_code), src_position_raw=VALUES(src_position_raw),
  210. employment_status=VALUES(employment_status), src_employment_status_raw=VALUES(src_employment_status_raw),
  211. written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  212. """,
  213. new { tid = tenantId, batch = batchId });
  214. }
  215. private async Task ProjectSalesLinesAsync(long tenantId, string batchId, CancellationToken cancellationToken)
  216. {
  217. if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_SALES", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
  218. return;
  219. var domain = await _domains.ResolveDomainAsync(MdpSourceIdentity.Native, tenantId, cancellationToken);
  220. await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_line");
  221. cancellationToken.ThrowIfCancellationRequested();
  222. await _db.Ado.ExecuteCommandAsync(
  223. $"""
  224. INSERT INTO mdp_std_work_order_line
  225. (tenant_id, factory_id, source_system, written_by, domain, doc_type, src_doc_type_raw,
  226. order_no, line_no, task_no, item_code, qty_planned, qty_completed, plan_finish_date, release_time,
  227. closed_flag, closed_time, approved_flag, void_flag, source_row_id, source_biz_key, sync_batch_id, sync_time)
  228. SELECT
  229. s.tenant_id, IFNULL(s.factory_id, 1), 'AIDOP_NATIVE', 'PLATFORM_FORM', @domain, 'SALES_ORDER', IFNULL(s.source_table,'mdp_std_so'),
  230. s.order_no, s.order_line, s.order_no, s.item_code, s.order_qty,
  231. COALESCE(d.shipped_qty, s.delivered_qty),
  232. COALESCE(s.promised_delivery_date, s.plan_delivery_date, s.customer_request_date), s.order_date,
  233. IF(IFNULL(s.closed,0)=1 OR IFNULL(s.line_closed_flag,0)=1
  234. OR d.shipped_qty>=s.order_qty, 1, 0),
  235. COALESCE(s.line_closed_time, IF(d.shipped_qty>=s.order_qty, d.last_ship_time, NULL)),
  236. 1, IFNULL(s.deleted_flag, 0),
  237. IFNULL(s.source_row_id, CAST(s.id AS CHAR)),
  238. LEFT(CONCAT('SO:', IFNULL(s.source_biz_key, s.order_no)), 200), @batch, NOW()
  239. FROM mdp_std_so s
  240. LEFT JOIN ({ShipmentActualsSql}) d
  241. ON d.tenant_id=s.tenant_id AND d.ref_task_no=s.order_no AND d.item_num=s.item_code
  242. WHERE s.tenant_id=@tid AND s.order_no IS NOT NULL AND s.order_no<>''
  243. AND s.source_system IN ('AIDOP_NATIVE','AIDOP')
  244. ON DUPLICATE KEY UPDATE
  245. item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), qty_completed=VALUES(qty_completed),
  246. plan_finish_date=VALUES(plan_finish_date), release_time=VALUES(release_time),
  247. closed_flag=VALUES(closed_flag), closed_time=VALUES(closed_time), void_flag=VALUES(void_flag),
  248. written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  249. """,
  250. new { tid = tenantId, batch = batchId, domain });
  251. }
  252. /// <summary>
  253. /// 用销售发货实绩刷新已有 <c>SALES_ORDER</c> 中立行的已交货量与关行状态。
  254. /// <para>
  255. /// 发货流水在 S5/S7 重算时才产出,而销售订单行投影只在 S1/S2 重算时跑;
  256. /// 若不在产出发货流水的同一次跑批里回写,S7 订单发货周期要等下一次 S1 重算才有值。
  257. /// 故本方法由 <c>InventoryMdpSyncService</c> 在发货投影之后调用,只更新既有行、不新增行。
  258. /// </para>
  259. /// </summary>
  260. public async Task RefreshSalesLineCompletionAsync(
  261. long tenantId, string? batchId, CancellationToken cancellationToken = default)
  262. {
  263. if (tenantId <= 0) return;
  264. if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_SALES", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
  265. return;
  266. cancellationToken.ThrowIfCancellationRequested();
  267. await _db.Ado.ExecuteCommandAsync(
  268. $"""
  269. UPDATE mdp_std_work_order_line l
  270. INNER JOIN ({ShipmentActualsSql}) d
  271. ON d.tenant_id=l.tenant_id AND d.ref_task_no=l.task_no AND d.item_num=l.item_code
  272. SET l.qty_completed=d.shipped_qty,
  273. l.closed_flag=IF(l.closed_flag=1 OR d.shipped_qty>=l.qty_planned, 1, 0),
  274. l.closed_time=COALESCE(l.closed_time, IF(d.shipped_qty>=l.qty_planned, d.last_ship_time, NULL)),
  275. l.sync_time=NOW()
  276. WHERE l.tenant_id=@tid AND l.doc_type='SALES_ORDER' AND l.source_system='AIDOP_NATIVE'
  277. """,
  278. new { tid = tenantId });
  279. }
  280. }