using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform; /// /// 把 T8 形态标准层投影到中立标准层。厂商枚举只允许出现在本文件(与主服务同属 T8 适配器)。 /// 指标 SQL 只读中立层。 /// public sealed partial class T8BaseInboundMdpSyncService { private const string NeutralBalanceConfigId = "t8_v5"; public async Task ProjectNeutralAsync( long tenantId, string batchId, DateTime now, CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); var ready = await _db.Ado.GetIntAsync( """ SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='mdp_std_inv_trans' AND COLUMN_NAME='src_trans_type_raw' """); if (ready == 0) return "中立层未就绪,已跳过投影"; var ps = new[] { new SugarParameter("@tid", tenantId), new SugarParameter("@batch", batchId), new SugarParameter("@now", now) }; var steps = new (string Name, string Sql)[] { ("inv_trans", ProjectInvTransSql), ("api_inv_trans", ProjectApiInvTransSql), ("api_inv_trans_unmapped", ProjectApiInvTransUnmappedSql), ("work_order_schedule", ProjectWorkOrderScheduleSql), ("work_order_line", ProjectWorkOrderLineSql), ("work_order_bom", ProjectWorkOrderBomSql), ("employee", ProjectEmployeeSql), ("fqc", ProjectFqcSql), ("report_start", ProjectReportStartSql) }; var failed = new List(); foreach (var step in steps) { cancellationToken.ThrowIfCancellationRequested(); try { await _db.Ado.ExecuteCommandAsync(step.Sql, ps); } catch (Exception ex) when (ex is not OperationCanceledException) { failed.Add($"{step.Name}: {ex.Message}"); } } return failed.Count == 0 ? "OK" : string.Join(" | ", failed); } /// /// 表值函数结果落入月度库存金额表。T8 未配置或调用失败时返回 0,不抛出,避免挡住推送现场。 /// public async Task TryMaterializeInventoryBalanceAsync( long tenantId, long factoryId, string domainCode, string periodYm, string batchId, DateTime now) { const string sqlTvf = @" select ckcode as ckcode, ckname as ckname, code as code, cname as cname, pcode as pcode, pname as pname, je3 as je3, je2 as je2 from dbo.Rep_总账_存货_V3(@domainCode, N'普通', N'正常', @startYm, @endYm)"; List rows; try { var t8 = _db.AsTenant().GetConnectionScope(NeutralBalanceConfigId); t8.Ado.CommandTimeOut = 60; rows = await t8.Ado.SqlQueryAsync(sqlTvf, new[] { new SugarParameter("@domainCode", domainCode), new SugarParameter("@startYm", periodYm), new SugarParameter("@endYm", periodYm) }); } catch (Exception ex) { Console.WriteLine($"[NeutralBalance] TVF skipped: {ex.Message}"); return 0; } var affected = 0; foreach (var r in rows) { var category = r.pcode ?? ""; var warehouse = r.ckcode ?? ""; var item = r.code ?? ""; affected += await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_inventory_balance_monthly (tenant_id, factory_id, source_system, domain, period_ym, category_code, category_name, warehouse_code, warehouse_name, item_code, avg_balance_amount, issue_cost_amount, source_biz_key, sync_batch_id, sync_time) VALUES (@tid, @factoryId, 'T8', @domain, @period, @category, @categoryName, @warehouse, @warehouseName, @item, @avgAmt, @issueAmt, @bizKey, @batch, @now) ON DUPLICATE KEY UPDATE category_name=VALUES(category_name), warehouse_name=VALUES(warehouse_name), avg_balance_amount=VALUES(avg_balance_amount), issue_cost_amount=VALUES(issue_cost_amount), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time) """, new SugarParameter("@tid", tenantId), new SugarParameter("@factoryId", factoryId), new SugarParameter("@domain", domainCode), new SugarParameter("@period", periodYm), new SugarParameter("@category", category), new SugarParameter("@categoryName", r.pname), new SugarParameter("@warehouse", warehouse), new SugarParameter("@warehouseName", r.ckname), new SugarParameter("@item", item), new SugarParameter("@avgAmt", r.je3 ?? 0m), new SugarParameter("@issueAmt", r.je2 ?? 0m), new SugarParameter("@bizKey", $"{domainCode}:{periodYm}:{category}:{warehouse}:{item}"), new SugarParameter("@batch", batchId), new SugarParameter("@now", now)); } return affected; } private const string BizDocCase = """ CASE h.lbs WHEN '采购入库' THEN 'PUR_RECEIPT' WHEN '生产领料' THEN 'PROD_ISSUE' WHEN '生产入库' THEN 'PROD_RECEIPT' WHEN '销售出库' THEN 'SALES_SHIP' WHEN '销售退货' THEN 'SALES_RETURN' ELSE 'OTHER' END """; private static readonly string ProjectInvTransSql = $@" INSERT INTO mdp_std_inv_trans (tenant_id, source_system, domain, src_rec_id, item_num, qty_change, doc_qty, trans_time, eff_date, approved_time, approved_flag, void_flag, summary_flag, biz_doc_type, src_biz_doc_type_raw, trans_type, ref_task_no, line_closed_flag, line_closed_time, history_from, as_of, sync_batch_id, dimension1, dimension2) SELECT l.tenant_id, 'T8', IFNULL(h.ztid,''), CAST(l.src_id AS CHAR), l.code, IFNULL(l.slzx,0), l.sl, l.addtime, h.date0, IFNULL(h.shtime, l.addtime), IFNULL(h.shyn,1), IFNULL(h.zfyn,0), IFNULL(h.hzyn,0), {BizDocCase}, h.lbs, CASE h.lbs WHEN '采购入库' THEN 'MAT_RECEIPT' WHEN '生产领料' THEN 'MAT_ISSUE' WHEN '生产入库' THEN 'FG_PROD_RECEIPT' WHEN '销售出库' THEN 'FG_SHIP' ELSE NULL END, l.lynoid, IFNULL(l.gdyn,0), l.gdtime, IFNULL(IFNULL(h.shtime, l.addtime), @now), @now, @batch, '', '' FROM mdp_std_t8_kc_tz_list l INNER JOIN mdp_std_t8_kc_tz_head h ON h.src_id=l.idid AND h.tenant_id=l.tenant_id WHERE l.tenant_id=@tid ON DUPLICATE KEY UPDATE item_num=VALUES(item_num), qty_change=VALUES(qty_change), doc_qty=VALUES(doc_qty), trans_time=VALUES(trans_time), eff_date=VALUES(eff_date), approved_time=VALUES(approved_time), approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag), summary_flag=VALUES(summary_flag), biz_doc_type=VALUES(biz_doc_type), src_biz_doc_type_raw=VALUES(src_biz_doc_type_raw), trans_type=VALUES(trans_type), ref_task_no=VALUES(ref_task_no), line_closed_flag=VALUES(line_closed_flag), line_closed_time=VALUES(line_closed_time), sync_batch_id=VALUES(sync_batch_id), as_of=VALUES(as_of)"; private static readonly string ProjectApiInvTransSql = BuildProjectApiInvTransSql(); private static string BuildProjectApiInvTransSql() { const string biz = "NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BizDocType')),'null')"; const string trans = "NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.TransType')),'null')"; const string src = "COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SrcTransType')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BizDocType')),'null'))"; return $@" INSERT INTO mdp_std_inv_trans (tenant_id, source_system, written_by, domain, src_rec_id, item_num, qty_change, location, trans_time, eff_date, approved_time, approved_flag, void_flag, summary_flag, biz_doc_type, src_biz_doc_type_raw, trans_type, src_trans_type_raw, ref_task_no, history_from, as_of, sync_batch_id, dimension1, dimension2) SELECT s.tenant_id, IFNULL(NULLIF(s.source_system,''),''), 'API_INBOUND', IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''), s.source_row_id, COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemCode')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.code')),'null')), CAST(COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Qty')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.sl')),'null')) AS DECIMAL(18,6)), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Location')),'null'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.TransTime')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.addtime')),'null')),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DocDate')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.date0')),'null')),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), IFNULL( STR_TO_DATE(SUBSTRING_INDEX(REPLACE(COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ApprovedTime')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.shtime')),'null')),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), STR_TO_DATE(SUBSTRING_INDEX(REPLACE(COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.TransTime')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.addtime')),'null')),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s')), IFNULL(CAST(COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ApprovedFlag')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.shyn')),'null')) AS SIGNED), 1), IFNULL(CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.VoidFlag')),'null') AS SIGNED), 0), IFNULL(CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SummaryFlag')),'null') AS SIGNED), 0), {NeutralTransTypeCodes.KnownBizDocCase(biz)}, {biz}, {NeutralTransTypeCodes.ApiTransTypeCase(trans, biz)}, {src}, COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.RefTaskNo')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.rwnoid')),'null')), IFNULL( STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ApprovedTime')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s'), @now), @now, @batch, '', '' FROM mdp_stg_t8_kc_tz_list s WHERE s.tenant_id=@tid AND s.source_row_id IS NOT NULL AND s.source_row_id<>'' ON DUPLICATE KEY UPDATE item_num=VALUES(item_num), qty_change=VALUES(qty_change), location=VALUES(location), approved_time=VALUES(approved_time), approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag), summary_flag=VALUES(summary_flag), biz_doc_type=VALUES(biz_doc_type), trans_type=VALUES(trans_type), src_trans_type_raw=VALUES(src_trans_type_raw), ref_task_no=VALUES(ref_task_no), sync_batch_id=VALUES(sync_batch_id)"; } private const string ProjectApiInvTransUnmappedSql = @" INSERT INTO mdp_std_inv_trans_unmapped (tenant_id, source_system, src_rec_id, src_trans_type_raw, unmapped_reason, location, row_count, first_seen, last_seen) SELECT tenant_id, source_system, src_rec_id, src_trans_type_raw, 'UNKNOWN_CODE', IFNULL(location,''), 1, NOW(), NOW() FROM mdp_std_inv_trans WHERE tenant_id=@tid AND sync_batch_id=@batch AND trans_type IS NULL AND src_trans_type_raw IS NOT NULL ON DUPLICATE KEY UPDATE last_seen=NOW(), row_count=row_count+1, unmapped_reason=VALUES(unmapped_reason), location=VALUES(location)"; private const string ProjectWorkOrderScheduleSql = @" INSERT INTO mdp_std_work_order_schedule (tenant_id, factory_id, source_system, work_order, doc_type, approved_flag, void_flag, source_biz_key, sync_batch_id, sync_time) SELECT h.tenant_id, 1, 'T8', h.noid, CASE h.lbs WHEN '生产任务' THEN 'PROD_TASK' WHEN '销售订单' THEN 'SALES_ORDER' ELSE 'OTHER' END, IFNULL(h.shyn,1), IFNULL(h.zf,0), CONCAT(IFNULL(h.ztid,''), ':', h.src_id), @batch, @now FROM mdp_std_t8_kc_dd_head h WHERE h.tenant_id=@tid AND h.noid IS NOT NULL AND h.noid<>'' AND NOT EXISTS ( SELECT 1 FROM mdp_std_work_order_schedule w WHERE w.tenant_id=h.tenant_id AND w.work_order=h.noid AND w.source_system<>'T8' AND w.doc_type = CASE h.lbs WHEN '销售订单' THEN 'SALES_ORDER' ELSE 'PROD_TASK' END) ON DUPLICATE KEY UPDATE approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag), source_biz_key=VALUES(source_biz_key), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)"; private const string ProjectWorkOrderLineSql = @" INSERT INTO mdp_std_work_order_line (tenant_id, factory_id, source_system, domain, doc_type, src_doc_type_raw, order_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 l.tenant_id, 1, 'T8', h.ztid, CASE h.lbs WHEN '生产任务' THEN 'PROD_TASK' WHEN '销售订单' THEN 'SALES_ORDER' ELSE 'OTHER' END, h.lbs, h.noid, l.rwnoid, l.code, l.sl, l.slzx, l.jhdate, l.addtime, IFNULL(l.gdyn,0), l.gdtime, IFNULL(h.shyn,1), IFNULL(h.zf,0), CAST(l.src_id AS CHAR), CONCAT(IFNULL(h.ztid,''), ':', l.src_id), @batch, @now FROM mdp_std_t8_kc_dd_list l INNER JOIN mdp_std_t8_kc_dd_head h ON h.src_id=l.idid AND h.tenant_id=l.tenant_id WHERE l.tenant_id=@tid AND h.noid IS NOT NULL AND h.noid<>'' ON DUPLICATE KEY UPDATE doc_type=VALUES(doc_type), task_no=VALUES(task_no), 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), approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)"; private const string ProjectWorkOrderBomSql = @" INSERT INTO mdp_std_work_order_bom (tenant_id, factory_id, source_system, domain, order_no, item_code, qty_required, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT l.tenant_id, 1, 'T8', h.ztid, h.noid, l.code, IFNULL(l.sl, 0), CAST(l.src_id AS CHAR), CONCAT(IFNULL(h.ztid,''), ':', l.src_id), @batch, @now FROM mdp_std_t8_kc_dd_list_cllist l INNER JOIN mdp_std_t8_kc_dd_head h ON h.src_id=l.idid AND h.tenant_id=l.tenant_id WHERE l.tenant_id=@tid AND h.noid IS NOT NULL AND h.noid<>'' ON DUPLICATE KEY UPDATE order_no=VALUES(order_no), item_code=VALUES(item_code), qty_required=VALUES(qty_required), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)"; private const string ProjectEmployeeSql = @" INSERT INTO mdp_std_employee (tenant_id, factory_id, source_system, domain, employee_no, position_code, src_position_raw, employment_status, src_employment_status_raw, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT p.tenant_id, 1, 'T8', IFNULL(p.ztid,''), CAST(p.src_id AS CHAR), CASE p.gw WHEN '计划' THEN 'PLANNER' WHEN '仓管' THEN 'WAREHOUSE' WHEN '生产' THEN 'PRODUCTION' ELSE 'OTHER' END, p.gw, CASE WHEN p.zzzt='在职' THEN 'ACTIVE' WHEN p.zzzt LIKE '%离职%' THEN 'LEFT' ELSE 'INACTIVE' END, p.zzzt, CAST(p.src_id AS CHAR), CONCAT(IFNULL(p.ztid,''), ':', p.src_id), @batch, @now FROM mdp_std_t8_sys_pelist p WHERE p.tenant_id=@tid AND p.src_id IS NOT NULL ON DUPLICATE KEY UPDATE 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), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)"; private const string ProjectFqcSql = @" INSERT INTO mdp_std_fqc_task (tenant_id, source_system, domain, apply_time, apply_flag, sales_order_line, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT z.tenant_id, 'T8', z.ztid, z.shdate, IFNULL(z.zjyn,0), CAST(z.lyid AS CHAR), CAST(z.src_id AS CHAR), CONCAT(IFNULL(z.ztid,''), ':', z.src_id), @batch, @now FROM mdp_std_t8_kc_zj_list z WHERE z.tenant_id=@tid AND z.src_id IS NOT NULL ON DUPLICATE KEY UPDATE domain=VALUES(domain), apply_time=VALUES(apply_time), apply_flag=VALUES(apply_flag), sales_order_line=VALUES(sales_order_line), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)"; private const string ProjectReportStartSql = @" INSERT INTO mdp_std_s6_report (tenant_id, source_system, work_order_no, start_work_date, domain, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT r.tenant_id, 'T8', r.noid, r.kgdate, r.ztid, CAST(r.src_id AS CHAR), CONCAT(IFNULL(r.ztid,''), ':', r.src_id), @batch, @now FROM mdp_std_t8_cj_bg_head_rep r WHERE r.tenant_id=@tid AND r.noid IS NOT NULL AND r.noid<>'' AND r.src_id IS NOT NULL ON DUPLICATE KEY UPDATE start_work_date=VALUES(start_work_date), domain=VALUES(domain), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)"; } internal sealed class T8InventoryBalanceRow { public string? ckcode { get; set; } public string? ckname { get; set; } public string? code { get; set; } public string? cname { get; set; } public string? pcode { get; set; } public string? pname { get; set; } public decimal? je3 { get; set; } public decimal? je2 { get; set; } }