| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292 |
- using Admin.NET.Plugin.AiDOP.DataPlatform.Schema;
- using Admin.NET.Plugin.AiDOP.Infrastructure;
- using Microsoft.Extensions.Logging;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform;
- /// <summary>
- /// 自有业务表投影到中立层。只写登记为 AIDOP_NATIVE 的对象。
- /// </summary>
- public sealed class NativeNeutralProjectionService : ITransient
- {
- private readonly ISqlSugarClient _db;
- private readonly MdpNeutralSourceGate _neutralGate;
- private readonly EmployeePositionMapService _positions;
- private readonly SourceDomainTenantResolver _domains;
- private readonly ILogger _logger;
- public NativeNeutralProjectionService(
- ISqlSugarClient db,
- MdpNeutralSourceGate neutralGate,
- EmployeePositionMapService positions,
- SourceDomainTenantResolver domains,
- ILoggerFactory loggerFactory)
- {
- _db = db;
- _neutralGate = neutralGate;
- _positions = positions;
- _domains = domains;
- _logger = loggerFactory.CreateLogger(nameof(NativeNeutralProjectionService));
- }
- /// <summary>
- /// 销售发货实绩(登记来源的 <c>SALES_SHIP</c>)按「订单号 + 物料」汇总。
- /// <para>
- /// 模式一下销售订单的权威在自有 S1,执行事实由外部系统提供,故订单行的已交货量与关行状态
- /// 要用发货流水推导,不能只看 <c>mdp_std_so.delivered_qty</c>(自建单常年为空)。
- /// 来源不写死:内联 <c>mdp_tenant_std_source</c>,只读该租户登记为 INV_TRANS 权威的那一个来源。
- /// 粒度与 S7 指标 SQL 的 JOIN 一致(中立流水没有订单行号,最细只能到订单号 + 物料)。
- /// </para>
- /// </summary>
- private const string ShipmentActualsSql = """
- SELECT t.tenant_id, t.ref_task_no, t.item_num,
- SUM(ABS(IFNULL(t.qty_change,0))) AS shipped_qty,
- MAX(t.approved_time) AS last_ship_time
- FROM mdp_std_inv_trans t
- INNER JOIN mdp_tenant_std_source ts
- ON ts.tenant_id=t.tenant_id AND ts.std_object='INV_TRANS' AND ts.source_system=t.source_system
- WHERE t.tenant_id=@tid AND t.biz_doc_type='SALES_SHIP'
- AND t.summary_flag=0 AND t.void_flag=0 AND t.approved_flag=1
- GROUP BY t.tenant_id, t.ref_task_no, t.item_num
- """;
- public async Task ProjectAsync(long tenantId, string batchId, CancellationToken cancellationToken = default)
- {
- if (tenantId <= 0) return;
- var failed = new List<string>();
- await Try(failed, "WO_LINE_PROD", () => ProjectWorkOrdersAsync(tenantId, batchId, cancellationToken));
- await Try(failed, "WO_SCHEDULE", () => ProjectWorkOrderScheduleAsync(tenantId, batchId, cancellationToken));
- await Try(failed, "WO_BOM", () => ProjectBomAsync(tenantId, batchId, cancellationToken));
- await Try(failed, "EMPLOYEE", () => ProjectEmployeesAsync(tenantId, batchId, cancellationToken));
- await Try(failed, "WO_LINE_SALES", () => ProjectSalesLinesAsync(tenantId, batchId, cancellationToken));
- if (failed.Count > 0)
- _logger.LogWarning("自有中立投影未完成 tenant={Tenant} {Failed}", tenantId, string.Join(" | ", failed));
- }
- private async Task Try(List<string> failed, string name, Func<Task> action)
- {
- try
- {
- await action();
- }
- catch (Exception ex) when (ex is not OperationCanceledException)
- {
- failed.Add($"{name}: {ex.Message}");
- }
- }
- private async Task ProjectWorkOrdersAsync(long tenantId, string batchId, CancellationToken cancellationToken)
- {
- if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_PROD", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
- return;
- await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_line");
- cancellationToken.ThrowIfCancellationRequested();
- 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, qty_completed, plan_finish_date, release_time,
- closed_flag, closed_time, approved_flag, void_flag, source_row_id, source_biz_key, sync_batch_id, sync_time)
- SELECT
- w.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(w.Domain,''), 'PROD_TASK', 'WorkOrdMaster',
- w.WorkOrd, NULL, w.WorkOrd, w.ItemNum, w.QtyOrded, w.QtyCompleted, w.DueDate, w.ReleaseDate,
- IF(UPPER(IFNULL(w.Status,''))='C', 1, 0), NULL, 1, IF(IFNULL(w.IsActive,1)=0, 1, 0),
- CAST(w.RecID AS CHAR), LEFT(CONCAT(IFNULL(w.Domain,''), ':', w.WorkOrd), 200), @batch, NOW()
- FROM WorkOrdMaster w
- WHERE w.tenant_id=@tid AND w.WorkOrd IS NOT NULL AND w.WorkOrd<>''
- ON DUPLICATE KEY UPDATE
- item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), qty_completed=VALUES(qty_completed),
- plan_finish_date=VALUES(plan_finish_date), release_time=VALUES(release_time),
- closed_flag=VALUES(closed_flag), void_flag=VALUES(void_flag),
- written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- new { tid = tenantId, batch = batchId });
- }
- /// <summary>
- /// 自有生产工单 → <c>mdp_std_work_order_schedule</c>(<c>doc_type='PROD_TASK'</c>)。
- /// <para>
- /// S5_L1_002 物料齐套满足率的分母以工单排程头为驱动,且 BOM 按 <c>source_system</c> 关联,
- /// 故排程必须与 <see cref="ProjectBomAsync"/> 写同一个 <c>AIDOP_NATIVE</c>,否则分母恒为空、
- /// 指标只会算出 NO_DATA。
- /// </para>
- /// <para>
- /// 唯一键 <c>uk_std_wo_sched_type</c> 不含 <c>source_system</c>,同一工单号若已被别的来源占用,
- /// 直接 upsert 会把对方的行改成自有来源,故与 T8 投影一样加 NOT EXISTS 反向占用保护。
- /// </para>
- /// </summary>
- private async Task ProjectWorkOrderScheduleAsync(long tenantId, string batchId, CancellationToken cancellationToken)
- {
- if (!await _neutralGate.AllowsAsync(tenantId, "WO_SCHEDULE", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
- return;
- await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_schedule");
- cancellationToken.ThrowIfCancellationRequested();
- await _db.Ado.ExecuteCommandAsync(
- """
- INSERT INTO mdp_std_work_order_schedule
- (tenant_id, factory_id, source_system, written_by, work_order, doc_type,
- item_code, item_name, site_code, status, priority, urgent_flag,
- qty_ordered, qty_completed, order_date, due_date, release_date,
- prod_line, lot_serial, drawing_no, project, work_order_type, labor_variance,
- approved_flag, void_flag, source_biz_key, sync_batch_id, sync_time)
- SELECT
- w.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', w.WorkOrd, 'PROD_TASK',
- w.ItemNum, w.ItemName, w.Site, w.Status, w.Priority, IF(IFNULL(w.Urgent,0)<>0, 1, 0),
- w.QtyOrded, w.QtyCompleted, w.OrdDate, w.DueDate, w.ReleaseDate,
- w.ProdLine, w.Batch, w.Drawing, w.Project, w.Typed, w.LbrVar,
- 1, IF(IFNULL(w.IsActive,1)=0, 1, 0),
- LEFT(CONCAT(IFNULL(w.Domain,''), ':', w.WorkOrd), 200), @batch, NOW()
- FROM WorkOrdMaster w
- WHERE w.tenant_id=@tid AND w.WorkOrd IS NOT NULL AND w.WorkOrd<>''
- AND NOT EXISTS (
- SELECT 1 FROM mdp_std_work_order_schedule x
- WHERE x.tenant_id=w.tenant_id AND x.doc_type='PROD_TASK' AND x.work_order=w.WorkOrd
- AND x.source_system<>'AIDOP_NATIVE')
- ON DUPLICATE KEY UPDATE
- item_code=VALUES(item_code), item_name=VALUES(item_name), site_code=VALUES(site_code),
- status=VALUES(status), priority=VALUES(priority), urgent_flag=VALUES(urgent_flag),
- qty_ordered=VALUES(qty_ordered), qty_completed=VALUES(qty_completed),
- order_date=VALUES(order_date), due_date=VALUES(due_date), release_date=VALUES(release_date),
- prod_line=VALUES(prod_line), lot_serial=VALUES(lot_serial), drawing_no=VALUES(drawing_no),
- project=VALUES(project), work_order_type=VALUES(work_order_type),
- labor_variance=VALUES(labor_variance),
- approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag),
- written_by=VALUES(written_by), source_biz_key=VALUES(source_biz_key),
- sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- new { tid = tenantId, batch = batchId });
- }
- private async Task ProjectBomAsync(long tenantId, string batchId, CancellationToken cancellationToken)
- {
- if (!await _neutralGate.AllowsAsync(tenantId, "WO_BOM", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
- return;
- await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_bom");
- cancellationToken.ThrowIfCancellationRequested();
- 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, task_no, item_code,
- qty_required, unit, source_row_id, source_biz_key, sync_batch_id, sync_time)
- SELECT
- d.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(d.Domain,''), d.WorkOrd,
- CAST(d.LineNum AS CHAR), d.WorkOrd, d.ItemNum, d.QtyRequired, d.UM,
- CAST(d.RecID AS CHAR),
- LEFT(CONCAT(IFNULL(d.Domain,''), ':', d.WorkOrd, ':', d.LineNum, ':', IFNULL(d.ItemNum,'')), 200),
- @batch, NOW()
- FROM WorkOrdDetail d
- WHERE d.tenant_id=@tid AND d.WorkOrd IS NOT NULL AND d.WorkOrd<>''
- ON DUPLICATE KEY UPDATE
- item_code=VALUES(item_code), qty_required=VALUES(qty_required), unit=VALUES(unit),
- written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- new { tid = tenantId, batch = batchId });
- }
- private async Task ProjectEmployeesAsync(long tenantId, string batchId, CancellationToken cancellationToken)
- {
- if (!await _neutralGate.AllowsAsync(tenantId, "EMPLOYEE", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
- return;
- await _positions.SyncNativeAsync(tenantId, cancellationToken);
- await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_employee");
- cancellationToken.ThrowIfCancellationRequested();
- await _db.Ado.ExecuteCommandAsync(
- """
- INSERT INTO mdp_std_employee
- (tenant_id, factory_id, source_system, written_by, domain, employee_no, employee_name, department_code,
- position_code, src_position_raw, employment_status, src_employment_status_raw,
- source_row_id, source_biz_key, sync_batch_id, sync_time)
- SELECT
- e.tenant_id, 1, 'AIDOP_NATIVE', 'PLATFORM_FORM', IFNULL(e.Domain,''), e.Employee, e.Name, e.Department,
- IFNULL(m.position_code, 'UNKNOWN'),
- COALESCE(NULLIF(e.JobTitle,''), NULLIF(e.WorkCtr,'')),
- CASE
- WHEN e.DateTerminated IS NOT NULL AND e.DateTerminated <= CURDATE() THEN 'LEFT'
- WHEN e.EmploymentStatus IN ('在职') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('ACTIVE','ONJOB','ON_JOB','AC') THEN 'ACTIVE'
- WHEN e.EmploymentStatus IN ('离职','辞退') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('LEFT','LEAVE','RESIGNED','TERMINATED','QUIT','DC') THEN 'LEFT'
- WHEN e.EmploymentStatus IN ('停用') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('INACTIVE','DISABLED','SUSPENDED') THEN 'INACTIVE'
- ELSE 'UNKNOWN' END,
- e.EmploymentStatus,
- CAST(e.RecID AS CHAR), LEFT(CONCAT(IFNULL(e.Domain,''), ':', e.Employee), 200), @batch, NOW()
- FROM EmployeeMaster e
- LEFT JOIN mdp_employee_position_map m
- ON m.tenant_id=e.tenant_id AND m.source_system='AIDOP_NATIVE'
- AND m.domain=IFNULL(e.Domain,'')
- AND m.src_position_raw=COALESCE(NULLIF(e.JobTitle,''), NULLIF(e.WorkCtr,''))
- WHERE e.tenant_id=@tid AND e.Employee IS NOT NULL AND e.Employee<>''
- ON DUPLICATE KEY UPDATE
- employee_name=VALUES(employee_name), department_code=VALUES(department_code),
- position_code=VALUES(position_code), src_position_raw=VALUES(src_position_raw),
- employment_status=VALUES(employment_status), src_employment_status_raw=VALUES(src_employment_status_raw),
- written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- new { tid = tenantId, batch = batchId });
- }
- private async Task ProjectSalesLinesAsync(long tenantId, string batchId, CancellationToken cancellationToken)
- {
- if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_SALES", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
- return;
- var domain = await _domains.ResolveDomainAsync(MdpSourceIdentity.Native, tenantId, cancellationToken);
- await MdpSchemaAligner.EnsureWrittenByColumnAsync(_db, "mdp_std_work_order_line");
- cancellationToken.ThrowIfCancellationRequested();
- 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, qty_completed, plan_finish_date, release_time,
- closed_flag, closed_time, approved_flag, void_flag, source_row_id, source_biz_key, sync_batch_id, sync_time)
- SELECT
- s.tenant_id, IFNULL(s.factory_id, 1), 'AIDOP_NATIVE', 'PLATFORM_FORM', @domain, 'SALES_ORDER', IFNULL(s.source_table,'mdp_std_so'),
- s.order_no, s.order_line, s.order_no, s.item_code, s.order_qty,
- COALESCE(d.shipped_qty, s.delivered_qty),
- COALESCE(s.promised_delivery_date, s.plan_delivery_date, s.customer_request_date), s.order_date,
- IF(IFNULL(s.closed,0)=1 OR IFNULL(s.line_closed_flag,0)=1
- OR d.shipped_qty>=s.order_qty, 1, 0),
- COALESCE(s.line_closed_time, IF(d.shipped_qty>=s.order_qty, d.last_ship_time, NULL)),
- 1, IFNULL(s.deleted_flag, 0),
- IFNULL(s.source_row_id, CAST(s.id AS CHAR)),
- LEFT(CONCAT('SO:', IFNULL(s.source_biz_key, s.order_no)), 200), @batch, NOW()
- FROM mdp_std_so s
- LEFT JOIN ({ShipmentActualsSql}) d
- ON d.tenant_id=s.tenant_id AND d.ref_task_no=s.order_no AND d.item_num=s.item_code
- WHERE s.tenant_id=@tid AND s.order_no IS NOT NULL AND s.order_no<>''
- AND s.source_system IN ('AIDOP_NATIVE','AIDOP')
- ON DUPLICATE KEY UPDATE
- item_code=VALUES(item_code), qty_planned=VALUES(qty_planned), qty_completed=VALUES(qty_completed),
- plan_finish_date=VALUES(plan_finish_date), release_time=VALUES(release_time),
- closed_flag=VALUES(closed_flag), closed_time=VALUES(closed_time), void_flag=VALUES(void_flag),
- written_by=VALUES(written_by), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- new { tid = tenantId, batch = batchId, domain });
- }
- /// <summary>
- /// 用销售发货实绩刷新已有 <c>SALES_ORDER</c> 中立行的已交货量与关行状态。
- /// <para>
- /// 发货流水在 S5/S7 重算时才产出,而销售订单行投影只在 S1/S2 重算时跑;
- /// 若不在产出发货流水的同一次跑批里回写,S7 订单发货周期要等下一次 S1 重算才有值。
- /// 故本方法由 <c>InventoryMdpSyncService</c> 在发货投影之后调用,只更新既有行、不新增行。
- /// </para>
- /// </summary>
- public async Task RefreshSalesLineCompletionAsync(
- long tenantId, string? batchId, CancellationToken cancellationToken = default)
- {
- if (tenantId <= 0) return;
- if (!await _neutralGate.AllowsAsync(tenantId, "WO_LINE_SALES", MdpSourceIdentity.Native, syncBatchId: batchId, cancellationToken: cancellationToken))
- return;
- cancellationToken.ThrowIfCancellationRequested();
- await _db.Ado.ExecuteCommandAsync(
- $"""
- UPDATE mdp_std_work_order_line l
- INNER JOIN ({ShipmentActualsSql}) d
- ON d.tenant_id=l.tenant_id AND d.ref_task_no=l.task_no AND d.item_num=l.item_code
- SET l.qty_completed=d.shipped_qty,
- l.closed_flag=IF(l.closed_flag=1 OR d.shipped_qty>=l.qty_planned, 1, 0),
- l.closed_time=COALESCE(l.closed_time, IF(d.shipped_qty>=l.qty_planned, d.last_ship_time, NULL)),
- l.sync_time=NOW()
- WHERE l.tenant_id=@tid AND l.doc_type='SALES_ORDER' AND l.source_system='AIDOP_NATIVE'
- """,
- new { tid = tenantId });
- }
- }
|