using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
///
/// S5 生产入库单 数据中台只读同步转换服务(DOP 内部,扁平单表,独立于 S5MdpSyncTransformService 的 KPI 管线)。
///
/// 源:aidopdev NbrMaster(n) + NbrDetail(d),业务类型 Type='WOI'(生产入库);
/// 扁平 LIST join:n.Domain=d.Domain AND n.Nbr=d.Nbr(一行=一条 NbrDetail 明细)。
/// 维表 DepartmentMaster(p) / ItemMaster(i) / LocationMaster(lt,lf) 全 LEFT JOIN。
/// LocationMaster 无 Domain 列(DOP 重建:domain_code/location/descr),按 tenant_id+location join。
/// 链路(Phase 1):执行器抽 NbrMaster/NbrDetail → mdp_stg_production_receipt → mdp_std_production_receipt。
/// 双模式入站:mdp_entity=S7_NBR_MASTER/S7_NBR_DETAIL → Writer → stg;标准层读 stg;维表仍直连本库。
///
/// 约束:
/// - 只读源/贴源,仅写 mdp_stg_production_receipt / mdp_std_production_receipt;绝不写 NbrMaster/NbrDetail。
/// - 过滤 Type='WOI' AND IsActive=1。
/// - WOI 为 0 行时转换成功完成、处理数为 0,不报错。
///
public class ProductionReceiptMdpSyncService : ITransient
{
private const string JobCode = "S5_PRODUCTION_RECEIPT_MDP_SYNC";
private const string InboundEntityCode = "S7_NBR_MASTER";
private const string InboundDetailEntityCode = "S7_NBR_DETAIL";
private readonly ISqlSugarClient _db;
private readonly MdpSourcePullDispatcher _pullDispatcher;
public ProductionReceiptMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher)
{
_db = db;
_pullDispatcher = pullDispatcher;
}
/// 全量:本地 DB 执行器灌 stg → 标准层(读 stg)。
public async Task RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
{
cancellationToken.ThrowIfCancellationRequested();
await EnsureTablesAsync();
await EnsureStgTableAsync();
var now = DateTime.Now;
var batchId = $"S5_PROD_RCPT_FULL_{now:yyyyMMddHHmmss}";
var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
try
{
var pullCtx = new MdpPullContext
{
TenantId = 0,
FullRefresh = true,
TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND",
BatchId = $"{batchId}_PULL"
};
await PopulateStgAsync(pullCtx, cancellationToken);
result.Rows = await TransformStandardAsync(batchId, now);
await MarkRunSuccessAsync(runLogId, now, result);
return result;
}
catch (Exception ex)
{
await MarkRunFailedAsync(runLogId, now, ex.Message);
throw;
}
}
///
/// S7 成品入库双模式入站:执行器抽主/明细 → stg,再跑 WOI 标准层(读 stg)。
///
public async Task RunInboundAsync(
long tenantId = 0,
bool fullRefresh = false,
CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
await EnsureTablesAsync();
await EnsureStgTableAsync();
var now = DateTime.Now;
var pullCtx = new MdpPullContext
{
TenantId = tenantId,
FullRefresh = fullRefresh,
TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND",
BatchId = $"S7_NBR_IN_{now:yyyyMMddHHmmss}"
};
var pull = await PopulateStgAsync(pullCtx, cancellationToken);
var batchId = $"S7_NBR_STD_{now:yyyyMMddHHmmss}";
var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
try
{
result.Rows = await TransformStandardAsync(batchId, now);
await MarkRunSuccessAsync(runLogId, now, result);
}
catch (Exception ex)
{
await MarkRunFailedAsync(runLogId, now, ex.Message);
throw;
}
return new ProductionReceiptInboundResult
{
PullBatchId = pullCtx.BatchId,
RowsPulled = pull.RowsPulled,
RowsWrittenStg = pull.RowsWritten,
NewCursor = pull.NewCursor,
TransformBatchId = result.BatchId,
StdRows = result.Rows
};
}
private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
MdpPullContext pullCtx, CancellationToken cancellationToken)
{
var master = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, pullCtx, cancellationToken);
var detail = await _pullDispatcher.PullByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
return (
master.RowsPulled + detail.RowsPulled,
master.RowsWritten + detail.RowsWritten,
detail.NewCursor ?? master.NewCursor);
}
private async Task EnsureStgTableAsync()
{
await _db.Ado.ExecuteCommandAsync(
"""
CREATE TABLE IF NOT EXISTS mdp_stg_production_receipt (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
tenant_id BIGINT NOT NULL,
source_system VARCHAR(50) NULL,
source_table VARCHAR(200),
source_row_id VARCHAR(200),
source_biz_key VARCHAR(300) NULL,
raw_data JSON,
sync_batch_id VARCHAR(100),
sync_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP,
process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
process_message VARCHAR(500) NULL,
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
KEY idx_batch (sync_batch_id),
KEY idx_src (source_table, source_row_id),
KEY idx_tenant (tenant_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='S7成品入库执行器贴源层'
""");
}
/// 防御式建表(与 UpdateScripts/1.0.214.sql 同构,幂等)。
private async Task EnsureTablesAsync()
{
await _db.Ado.ExecuteCommandAsync(
"""
CREATE TABLE IF NOT EXISTS mdp_std_production_receipt (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
tenant_id BIGINT NOT NULL DEFAULT 0,
factory_id BIGINT NULL DEFAULT 1,
source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
domain VARCHAR(80) NOT NULL,
master_rec_id INT NULL,
detail_rec_id INT NOT NULL,
nbr VARCHAR(24) NULL,
line SMALLINT NOT NULL DEFAULT 0,
receipt_date DATETIME NULL,
status VARCHAR(8) NULL,
status_desc VARCHAR(20) NULL,
remark VARCHAR(200) NULL,
prod_line VARCHAR(8) NULL,
work_ord VARCHAR(64) NULL,
erp_work_ord VARCHAR(60) NULL,
department VARCHAR(8) NULL,
department_desc VARCHAR(255) NULL,
applicant_name VARCHAR(12) NULL,
item_num VARCHAR(24) NULL,
item_name VARCHAR(200) NULL,
item_spec VARCHAR(200) NULL,
um VARCHAR(8) NULL,
location_to VARCHAR(8) NULL,
location_to_desc VARCHAR(255) NULL,
lot_serial VARCHAR(120) NULL,
qty_rec DECIMAL(18,5) NULL DEFAULT 0,
qty_to DECIMAL(18,5) NULL DEFAULT 0,
location_from VARCHAR(8) NULL,
location_from_desc VARCHAR(255) NULL,
ord_nbr VARCHAR(48) NULL,
source_biz_key VARCHAR(200) NULL,
sync_batch_id VARCHAR(100) NOT NULL,
sync_time DATETIME NOT NULL,
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_mdp_std_prod_receipt (tenant_id, domain, detail_rec_id),
KEY idx_mdp_std_prod_rcpt_nbr (tenant_id, nbr),
KEY idx_mdp_std_prod_rcpt_date (tenant_id, receipt_date),
KEY idx_mdp_std_prod_rcpt_workord (tenant_id, work_ord),
KEY idx_mdp_std_prod_rcpt_erp (tenant_id, erp_work_ord),
KEY idx_mdp_std_prod_rcpt_lot (tenant_id, lot_serial),
KEY idx_mdp_std_prod_rcpt_item (tenant_id, item_num),
KEY idx_mdp_std_prod_rcpt_locto (tenant_id, location_to)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S5生产入库单标准层(扁平,一行=一条明细)'
""");
}
///
/// 标准化:stg(NbrMaster/NbrDetail) + 维表(本库) → mdp_std_production_receipt。
///
private async Task TransformStandardAsync(string batchId, DateTime now)
{
var rows = await _db.Ado.GetIntAsync(
"""
SELECT COUNT(1)
FROM mdp_stg_production_receipt n
INNER JOIN mdp_stg_production_receipt d
ON d.source_table='NbrDetail'
AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
WHERE n.source_table='NbrMaster'
AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1
AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL
""");
await _db.Ado.ExecuteCommandAsync(
"""
INSERT INTO mdp_std_production_receipt
(tenant_id, factory_id, source_system, domain, master_rec_id, detail_rec_id, nbr, line,
receipt_date, status, status_desc, remark, prod_line, work_ord, erp_work_ord,
department, department_desc, applicant_name, item_num, item_name, item_spec, um,
location_to, location_to_desc, lot_serial, qty_rec, qty_to, location_from, location_from_desc,
ord_nbr, source_biz_key, sync_batch_id, sync_time)
SELECT
IFNULL(n.tenant_id, 0), 1, IFNULL(NULLIF(n.source_system,''), 'AIDOP'),
JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')),
CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.RecID')) AS SIGNED),
CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS SIGNED),
JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')),
CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS SIGNED),
STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Date')),'null'),''), '%Y-%m-%d %H:%i:%s'),
UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Remark')),
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ProdLine')),
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.WorkOrd')),
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Address')),
JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),
TRIM(CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),''), ' ', IFNULL(p.Descr, ''))),
JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Name')),
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')), i.Descr, i.Descr1, i.UM,
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')), lt.descr,
CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LotSerial')) AS CHAR),
CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyRec')) AS DECIMAL(18,5)),
CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyTo')) AS DECIMAL(18,5)),
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')), lf.descr,
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')),
CONCAT(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), '#',
JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID'))),
@BatchId, @Now
FROM mdp_stg_production_receipt n
INNER JOIN mdp_stg_production_receipt d
ON d.source_table='NbrDetail'
AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
LEFT JOIN DepartmentMaster p
ON p.Domain = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain'))
AND p.Department = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department'))
LEFT JOIN ItemMaster i
ON i.Domain = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
AND i.ItemNum = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum'))
LEFT JOIN LocationMaster lt
ON lt.tenant_id = IFNULL(d.tenant_id, 0)
AND lt.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo'))
LEFT JOIN LocationMaster lf
ON lf.tenant_id = IFNULL(d.tenant_id, 0)
AND lf.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom'))
WHERE n.source_table='NbrMaster'
AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1
AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL
ON DUPLICATE KEY UPDATE
factory_id=VALUES(factory_id), master_rec_id=VALUES(master_rec_id), nbr=VALUES(nbr), line=VALUES(line),
receipt_date=VALUES(receipt_date), status=VALUES(status), status_desc=VALUES(status_desc), remark=VALUES(remark),
prod_line=VALUES(prod_line), work_ord=VALUES(work_ord), erp_work_ord=VALUES(erp_work_ord),
department=VALUES(department), department_desc=VALUES(department_desc), applicant_name=VALUES(applicant_name),
item_num=VALUES(item_num), item_name=VALUES(item_name), item_spec=VALUES(item_spec), um=VALUES(um),
location_to=VALUES(location_to), location_to_desc=VALUES(location_to_desc), lot_serial=VALUES(lot_serial),
qty_rec=VALUES(qty_rec), qty_to=VALUES(qty_to), location_from=VALUES(location_from),
location_from_desc=VALUES(location_from_desc), ord_nbr=VALUES(ord_nbr),
sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
""",
new SugarParameter("@BatchId", batchId),
new SugarParameter("@Now", now));
return rows;
}
private async Task InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
{
await _db.Ado.ExecuteCommandAsync(
"""
INSERT INTO mdp_transform_run_log
(tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
VALUES (0, @JobCode, 'S5生产入库单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
""",
new SugarParameter("@JobCode", JobCode),
new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
new SugarParameter("@BatchId", batchId),
new SugarParameter("@StartTime", startedAt));
return await _db.Ado.GetLongAsync(
"SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
new List { new("@BatchId", batchId) });
}
private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, ProductionReceiptMdpSyncResult result)
{
var finishedAt = DateTime.Now;
await _db.Ado.ExecuteCommandAsync(
"""
UPDATE mdp_transform_run_log
SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
WHERE id=@Id
""",
new SugarParameter("@EndTime", finishedAt),
new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
new SugarParameter("@StandardRows", result.Rows),
new SugarParameter("@Id", runLogId));
}
private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
{
var finishedAt = DateTime.Now;
await _db.Ado.ExecuteCommandAsync(
"""
UPDATE mdp_transform_run_log
SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
WHERE id=@Id
""",
new SugarParameter("@EndTime", finishedAt),
new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
new SugarParameter("@Id", runLogId));
}
private static string NormalizeTriggerType(string? triggerType)
=> string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
}
/// 生产入库单 MDP 同步转换结果。
public sealed class ProductionReceiptMdpSyncResult
{
public long RunLogId { get; set; }
public string BatchId { get; set; } = string.Empty;
public int Rows { get; set; }
}
/// S7 成品入库双模式入站结果。
public sealed class ProductionReceiptInboundResult
{
public string PullBatchId { get; set; } = string.Empty;
public int RowsPulled { get; set; }
public int RowsWrittenStg { get; set; }
public string? NewCursor { get; set; }
public string TransformBatchId { get; set; } = string.Empty;
public int StdRows { get; set; }
}