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') THEN 'ACTIVE'
  196. WHEN e.EmploymentStatus IN ('离职','辞退') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('LEFT','LEAVE','RESIGNED','TERMINATED','QUIT') 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. }