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; }
}