using Admin.NET.Plugin.AiDOP.DataPlatform.Schema;
using Admin.NET.Plugin.AiDOP.Infrastructure;
using Microsoft.Extensions.Logging;
namespace Admin.NET.Plugin.AiDOP.DataPlatform;
///
/// 自有业务表投影到中立层。只写登记为 AIDOP_NATIVE 的对象。
///
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));
}
///
/// 销售发货实绩(登记来源的 SALES_SHIP)按「订单号 + 物料」汇总。
///
/// 模式一下销售订单的权威在自有 S1,执行事实由外部系统提供,故订单行的已交货量与关行状态
/// 要用发货流水推导,不能只看 mdp_std_so.delivered_qty(自建单常年为空)。
/// 来源不写死:内联 mdp_tenant_std_source,只读该租户登记为 INV_TRANS 权威的那一个来源。
/// 粒度与 S7 指标 SQL 的 JOIN 一致(中立流水没有订单行号,最细只能到订单号 + 物料)。
///
///
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();
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 failed, string name, Func 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 });
}
///
/// 自有生产工单 → mdp_std_work_order_schedule(doc_type='PROD_TASK')。
///
/// S5_L1_002 物料齐套满足率的分母以工单排程头为驱动,且 BOM 按 source_system 关联,
/// 故排程必须与 写同一个 AIDOP_NATIVE,否则分母恒为空、
/// 指标只会算出 NO_DATA。
///
///
/// 唯一键 uk_std_wo_sched_type 不含 source_system,同一工单号若已被别的来源占用,
/// 直接 upsert 会把对方的行改成自有来源,故与 T8 投影一样加 NOT EXISTS 反向占用保护。
///
///
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') THEN 'ACTIVE'
WHEN e.EmploymentStatus IN ('离职','辞退') OR UPPER(IFNULL(e.EmploymentStatus,'')) IN ('LEFT','LEAVE','RESIGNED','TERMINATED','QUIT') 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 });
}
///
/// 用销售发货实绩刷新已有 SALES_ORDER 中立行的已交货量与关行状态。
///
/// 发货流水在 S5/S7 重算时才产出,而销售订单行投影只在 S1/S2 重算时跑;
/// 若不在产出发货流水的同一次跑批里回写,S7 订单发货周期要等下一次 S1 重算才有值。
/// 故本方法由 InventoryMdpSyncService 在发货投影之后调用,只更新既有行、不新增行。
///
///
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 });
}
}