using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
using Microsoft.Extensions.Logging;
using SqlSugar;
namespace Admin.NET.Plugin.AiDOP.Order;
///
/// S8 Stage-1「订单评审」ActualHours 数据链:
/// aidop_action_run_log → mdp_stg_action_run → mdp_std_action_event → dwd_order_review_kpi。
///
/// 为什么必须经中台:S8 正式运行时禁止直读 aidop_action_run_log /
/// crm_seorder / crm_seorderentry。配对、FAILED 过滤、重复动作处理全部在本服务完成,
/// S8 只消费 dwd_order_review_kpi。
///
/// 租户是唯一隔离边界:每一层的读、写、JOIN、唯一键首列都带 tenant_id;
/// 源侧 tenant_id 为 NULL 或 0 一律拒绝入站,绝不补默认租户。
///
public sealed class ActionEventMdpSyncService : ITransient
{
/// MDP 实体码(注册见 UpdateScripts/1.0.527.sql)。
public const string EntityCode = "S1_ACTION_RUN_LOG";
private const string SourceTable = "aidop_action_run_log";
private const string ReviewActionCode = "S1_ORDER_REVIEW";
private const string ConfirmActionCode = "S1_DELIVERY_CONFIRM";
private const string OrderBizType = "crm_seorder";
private readonly ISqlSugarClient _db;
private readonly MdpModuleStagingPuller _stagingPuller;
private readonly ILogger _logger;
public ActionEventMdpSyncService(
ISqlSugarClient db,
MdpModuleStagingPuller stagingPuller,
ILogger logger)
{
_db = db;
_stagingPuller = stagingPuller;
_logger = logger;
}
///
/// 全链路跑批:源→STG(增量 + RUNNING 回捞)→ STD → DWD KPI。
/// 幂等:同一份源数据重复执行,STD 与 DWD 结果完全一致。
///
public async Task RunAsync(long tenantId, CancellationToken cancellationToken = default)
{
if (tenantId <= 0)
throw new InvalidOperationException("ActionEventMdpSyncService 必须显式传入有效 tenantId,禁止默认租户。");
var now = DateTime.Now;
var batchId = $"S1_ACTION_EVENT_{tenantId}_{now:yyyyMMddHHmmss}";
var result = new ActionEventSyncResult { BatchId = batchId, TenantId = tenantId };
result.StageRows = await PullStagingAsync(tenantId, batchId, cancellationToken);
result.StageRepulledRunning = await RepullRunningAsync(tenantId, batchId, now, cancellationToken);
result.StandardRows = await TransformStdFromStgAsync(tenantId, batchId, cancellationToken);
result.KpiRows = await BuildOrderReviewKpiAsync(tenantId, batchId, now, cancellationToken);
// Stage-2 必须在 Stage-1 KPI 之后跑:它直接消费 dwd_order_review_kpi.review_end_time
// 作为 Start,而不是自己再实现一遍 Stage-1 的配对逻辑。
result.ProductDesignKpiRows = await BuildProductDesignKpiAsync(tenantId, batchId, now, cancellationToken);
_logger.LogInformation(
"S1 动作事件链路完成 tenant={TenantId} batch={BatchId} stg={Stg} repull={Repull} std={Std} kpi={Kpi} pdKpi={PdKpi}",
tenantId, batchId, result.StageRows, result.StageRepulledRunning,
result.StandardRows, result.KpiRows, result.ProductDesignKpiRows);
return result;
}
// ══════════════════════════════════════════════════════════════════════
// PHASE 2 · 源 → STG
// ══════════════════════════════════════════════════════════════════════
///
/// 主增量:走统一执行器(sync_mode=INCR,主游标 start_time,水位存
/// mdp_entity.last_cursor)。写入经 的 canonical 契约,
/// 其中已内置「tenant 解析不出即跳过 + 跨租户行丢弃」双重守卫。
///
private Task PullStagingAsync(long tenantId, string batchId, CancellationToken cancellationToken) =>
_stagingPuller.PullEntitiesAsync(
[EntityCode], batchId, tenantId,
fullRefresh: false,
taskCode: "S1_ACTION_EVENT_INBOUND",
cancellationToken,
factoryId: 0,
requireMatchingSourceTenant: true);
///
/// RUNNING 回捞补偿(本批的核心风险控制点)。
///
/// 源行不是内容级 append-only:StartAsync 先 INSERT 一行
/// status='RUNNING', end_time=NULL,之后 FinishAsync /
/// FailStaleRunningAsync 会原地 UPDATE 成 SUCCESS/FAILED 并补 end_time。
/// 而源表没有 update_time 列,且 create_time ≡ start_time、不随 UPDATE 推进
/// —— 所以主游标 start_time > @cursor 永远不会第二次捞到这一行,
/// 若无补偿,STG/STD 会永久停留在 RUNNING。
///
/// 补偿判据:以标准层当前仍为 RUNNING 的 action_run_id 集合反向驱动重拉。
/// 该集合天然极小(实测全库仅 1 行),因此:
/// ① tenant-safe —— 读写两侧都带 tenant_id=@TenantId;
/// ② 幂等 —— ON DUPLICATE KEY UPDATE,重复执行不产生新行;
/// ③ 不全表扫描 —— 由 RUNNING 集合驱动,走 uk_std_action_event 与源主键;
/// ④ 不依赖 id 单调性 —— 按 id 精确等值匹配,与雪花乱序无关;
/// ⑤ 不遗漏迟到终态 —— 只要某行在 STD 还是 RUNNING,每轮都会被重新检查。
///
/// 此处不经 而用定向 INSERT…SELECT,
/// 原因是框架执行器只支持「按水位翻页」、不支持「按 id 集合定向重拉」;
/// 因此租户守卫在本 SQL 内显式重建(tenant_id=@TenantId +
/// tenant_id IS NOT NULL + tenant_id<>0),语义与框架守卫一致。
///
private async Task RepullRunningAsync(
long tenantId, string batchId, DateTime now, CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
var ps = new SugarParameter[]
{
new("@TenantId", tenantId),
new("@BatchId", batchId),
new("@Now", now)
};
// 逻辑计数:本轮被回捞的 distinct action_run_id 数。
// 必须先于 upsert 求值(upsert 不改变该集合,但顺序固定更易解释),
// 且与下方 INSERT 共用同一份 WHERE(RepullRunningFrom),杜绝谓词漂移。
var logicalRows = await _db.Ado.GetIntAsync(
$"SELECT COUNT(DISTINCT s.`id`) {RepullRunningFrom}", ps);
const string sql = $"""
INSERT INTO `mdp_stg_action_run`
(`tenant_id`, `source_system`, `source_table`, `source_row_id`, `source_biz_key`,
`raw_data`, `sync_batch_id`, `sync_time`, `process_status`, `create_time`, `update_time`)
SELECT
s.`tenant_id`,
'AIDOPDEV_MYSQL',
'aidop_action_run_log',
CAST(s.`id` AS CHAR),
CAST(s.`id` AS CHAR),
JSON_OBJECT(
'id', CAST(s.`id` AS CHAR),
'tenant_id', s.`tenant_id`,
'action_code', s.`action_code`,
'biz_type', s.`biz_type`,
'biz_id', s.`biz_id`,
'biz_no', s.`biz_no`,
'status', s.`status`,
'message', s.`message`,
'detail_json', s.`detail_json`,
'start_time', DATE_FORMAT(s.`start_time`, '%Y-%m-%d %H:%i:%s'),
'end_time', DATE_FORMAT(s.`end_time`, '%Y-%m-%d %H:%i:%s'),
'create_time', DATE_FORMAT(s.`create_time`,'%Y-%m-%d %H:%i:%s')
),
@BatchId, @Now, 'PENDING', @Now, @Now
{RepullRunningFrom}
ON DUPLICATE KEY UPDATE
`source_row_id` = VALUES(`source_row_id`),
`raw_data` = VALUES(`raw_data`),
`sync_batch_id` = VALUES(`sync_batch_id`),
`sync_time` = VALUES(`sync_time`),
`process_status` = 'PENDING',
`update_time` = VALUES(`update_time`)
""";
// 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
await _db.Ado.ExecuteCommandAsync(sql, ps);
if (logicalRows > 0)
_logger.LogInformation("RUNNING 回捞补偿命中 tenant={TenantId} rows={Rows}", tenantId, logicalRows);
return logicalRows;
}
///
/// RUNNING 回捞的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
///
private const string RepullRunningFrom = """
FROM `aidop_action_run_log` s
WHERE s.`tenant_id` = @TenantId
AND s.`tenant_id` IS NOT NULL
AND s.`tenant_id` <> 0
AND EXISTS (
SELECT 1 FROM `mdp_std_action_event` e
WHERE e.`tenant_id` = @TenantId
AND e.`action_status` = 'RUNNING'
AND e.`action_run_id` = s.`id`
)
""";
// ══════════════════════════════════════════════════════════════════════
// PHASE 3 · STG → STD(审计事实层,全状态保留)
// ══════════════════════════════════════════════════════════════════════
///
/// 贴源 → 标准事件层。Grain = tenant_id + action_run_id。
/// RUNNING / SUCCESS / FAILED 全部保留 —— STD 是审计事实层,
/// 不因为 KPI 只需要 SUCCESS 就在此丢弃其它状态。
/// 同一 action_run_id 重拉时,ON DUPLICATE KEY UPDATE 允许
/// RUNNING 被 SUCCESS/FAILED 更新(终态推进)。
///
public async Task TransformStdFromStgAsync(
long tenantId, string batchId, CancellationToken cancellationToken = default)
{
if (tenantId <= 0) throw new InvalidOperationException("TransformStdFromStgAsync 需要有效 tenantId。");
cancellationToken.ThrowIfCancellationRequested();
var ps = new SugarParameter[]
{
new("@TenantId", tenantId),
new("@BatchId", batchId)
};
// 逻辑计数:本轮参与 STD transform 的 distinct action_run_id 数。
// 与下方 INSERT 共用同一份 FROM/WHERE(StdEligibleFrom),保证「计了什么」= 「写了什么」。
var logicalRows = await _db.Ado.GetIntAsync(
$"SELECT COUNT(DISTINCT JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id'))) {StdEligibleFrom}", ps);
const string sql = $"""
INSERT INTO `mdp_std_action_event`
(`tenant_id`, `action_run_id`, `action_code`, `biz_type`, `biz_id`, `biz_no`,
`action_status`, `action_start_time`, `action_end_time`, `message`, `detail_json`,
`source_system`, `source_table`, `source_biz_key`, `sync_batch_id`)
SELECT
g.`tenant_id`,
CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) AS SIGNED),
JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')),
JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')),
CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) REGEXP '^-?[0-9]+$'
THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) AS SIGNED) END,
NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_no')), 'null'),
JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')),
CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'T', ' ') AS DATETIME),
CASE WHEN NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'null') IS NULL
THEN NULL
ELSE CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'T', ' ') AS DATETIME) END,
NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.message')), 'null'),
CASE
WHEN JSON_EXTRACT(g.`raw_data`, '$.detail_json') IS NULL THEN NULL
WHEN JSON_VALID(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')))
THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')) AS JSON)
ELSE JSON_EXTRACT(g.`raw_data`, '$.detail_json')
END,
g.`source_system`,
g.`source_table`,
g.`source_biz_key`,
@BatchId
{StdEligibleFrom}
ON DUPLICATE KEY UPDATE
`action_code` = VALUES(`action_code`),
`biz_type` = VALUES(`biz_type`),
`biz_id` = VALUES(`biz_id`),
`biz_no` = VALUES(`biz_no`),
`action_status` = VALUES(`action_status`),
`action_start_time` = VALUES(`action_start_time`),
`action_end_time` = VALUES(`action_end_time`),
`message` = VALUES(`message`),
`detail_json` = VALUES(`detail_json`),
`sync_batch_id` = VALUES(`sync_batch_id`)
""";
// 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
await _db.Ado.ExecuteCommandAsync(sql, ps);
return logicalRows;
}
///
/// STD transform 的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
///
private const string StdEligibleFrom = """
FROM `mdp_stg_action_run` g
WHERE g.`tenant_id` = @TenantId
AND g.`tenant_id` <> 0
AND g.`source_table` = 'aidop_action_run_log'
AND JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) REGEXP '^[0-9]+$'
AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')), 'null') IS NOT NULL
AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')), 'null') IS NOT NULL
AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')), 'null') IS NOT NULL
AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'null') IS NOT NULL
""";
// ══════════════════════════════════════════════════════════════════════
// PHASE 4 · STD → DWD KPI(FIRST VALID PAIR)
// ══════════════════════════════════════════════════════════════════════
///
/// 标准事件层 → 订单评审 KPI。Grain = tenant_id + sales_order_id。
///
/// 算法(定版,不得改用 LAST PAIR / MAX):
/// ① 只取 status='SUCCESS',FAILED 保留在 STD 但不参与配对;
/// ② ReviewStart = 评审 SUCCESS 的 FIRST(ORDER BY action_start_time, action_run_id);
/// ③ ReviewEnd = 首个 action_end_time > ReviewStart 的确认 SUCCESS 的
/// action_end_time(= 确认动作完成时刻,是 progress='3' 提交后的最近上界);
/// ④ ActualHours = 自然 elapsed 小时,2 位小数,不扣夜间/周末/节假日(本项目无工作日历 Authority)。
///
/// 冻结语义:算法只依赖「第一条评审」与「首个晚于它的确认」,两者一旦产生即恒定,
/// 后续任意次重复点击都不会改写已形成的 KPI;因此重复执行结果完全一致(deterministic / idempotent)。
///
/// 本批不计算 TargetHours / kpi_status —— 那依赖 S0 PI 匹配,属另一批次。
///
public async Task BuildOrderReviewKpiAsync(
long tenantId, string batchId, DateTime now, CancellationToken cancellationToken = default)
{
if (tenantId <= 0) throw new InvalidOperationException("BuildOrderReviewKpiAsync 需要有效 tenantId。");
cancellationToken.ThrowIfCancellationRequested();
var ps = new SugarParameter[]
{
new("@TenantId", tenantId),
new("@BatchId", batchId),
new("@Now", now)
};
// 逻辑计数:本轮参与 KPI transform 的 distinct (tenant_id, sales_order_id) 数,
// 即下方 CTE `agg` 的基数 —— 与写入集合一一对应(agg 是最终 SELECT 的驱动表)。
var logicalRows = await _db.Ado.GetIntAsync($"""
SELECT COUNT(*) FROM (
SELECT `tenant_id`, `biz_id` {KpiEventFrom} GROUP BY `tenant_id`, `biz_id`
) agg
""", ps);
const string sql = $"""
INSERT INTO `dwd_order_review_kpi`
(`tenant_id`, `sales_order_id`, `sales_order_no`,
`review_action_run_id`, `confirm_action_run_id`,
`review_start_time`, `review_end_time`, `actual_hours`, `calculation_status`,
`review_success_count`, `confirm_success_count`,
`review_failed_count`, `confirm_failed_count`,
`calc_batch_id`, `calc_time`)
WITH ev AS (
SELECT `tenant_id`, `biz_id`, `biz_no`, `action_run_id`, `action_code`,
`action_status`, `action_start_time`, `action_end_time`
{KpiEventFrom}
),
agg AS (
SELECT `tenant_id`, `biz_id`,
MAX(`biz_no`) AS biz_no,
SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS') AS rs,
SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'SUCCESS') AS cs,
SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'FAILED') AS rf,
SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'FAILED') AS cf
FROM ev GROUP BY `tenant_id`, `biz_id`
),
rv AS (
SELECT `tenant_id`, `biz_id`, `action_run_id`, `action_start_time`,
ROW_NUMBER() OVER (PARTITION BY `tenant_id`, `biz_id`
ORDER BY `action_start_time` ASC, `action_run_id` ASC) AS rn
FROM ev
WHERE `action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS'
),
r1 AS (SELECT * FROM rv WHERE rn = 1),
cv AS (
SELECT c.`tenant_id`, c.`biz_id`, c.`action_run_id`, c.`action_end_time`,
ROW_NUMBER() OVER (PARTITION BY c.`tenant_id`, c.`biz_id`
ORDER BY c.`action_end_time` ASC, c.`action_run_id` ASC) AS rn
FROM ev c
JOIN r1 ON r1.`tenant_id` = c.`tenant_id` AND r1.`biz_id` = c.`biz_id`
WHERE c.`action_code` = 'S1_DELIVERY_CONFIRM'
AND c.`action_status` = 'SUCCESS'
AND c.`action_end_time` IS NOT NULL
AND c.`action_end_time` > r1.`action_start_time`
),
c1 AS (SELECT * FROM cv WHERE rn = 1)
SELECT
a.`tenant_id`,
a.`biz_id`,
a.biz_no,
r1.`action_run_id`,
c1.`action_run_id`,
r1.`action_start_time`,
c1.`action_end_time`,
CASE WHEN r1.`action_start_time` IS NOT NULL AND c1.`action_end_time` IS NOT NULL
THEN ROUND(TIMESTAMPDIFF(SECOND, r1.`action_start_time`, c1.`action_end_time`) / 3600.0, 2)
END,
CASE
WHEN r1.`action_start_time` IS NULL AND a.cs = 0 THEN 'NO_START'
WHEN r1.`action_start_time` IS NULL THEN 'INVALID'
WHEN c1.`action_end_time` IS NOT NULL THEN 'OK'
WHEN a.cs = 0 THEN 'IN_PROGRESS'
ELSE 'INVALID_TIME_SEQUENCE'
END,
a.rs, a.cs, a.rf, a.cf,
@BatchId, @Now
FROM agg a
LEFT JOIN r1 ON r1.`tenant_id` = a.`tenant_id` AND r1.`biz_id` = a.`biz_id`
LEFT JOIN c1 ON c1.`tenant_id` = a.`tenant_id` AND c1.`biz_id` = a.`biz_id`
ON DUPLICATE KEY UPDATE
`sales_order_no` = VALUES(`sales_order_no`),
`review_action_run_id` = VALUES(`review_action_run_id`),
`confirm_action_run_id` = VALUES(`confirm_action_run_id`),
`review_start_time` = VALUES(`review_start_time`),
`review_end_time` = VALUES(`review_end_time`),
`actual_hours` = VALUES(`actual_hours`),
`calculation_status` = VALUES(`calculation_status`),
`review_success_count` = VALUES(`review_success_count`),
`confirm_success_count` = VALUES(`confirm_success_count`),
`review_failed_count` = VALUES(`review_failed_count`),
`confirm_failed_count` = VALUES(`confirm_failed_count`),
`calc_batch_id` = VALUES(`calc_batch_id`),
`calc_time` = VALUES(`calc_time`)
""";
// 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
await _db.Ado.ExecuteCommandAsync(sql, ps);
return logicalRows;
}
///
/// KPI transform 的事件取数范围(FROM + WHERE)。被计数查询与 CTE ev 共用,
/// 保证「计的订单数」与「写入的订单集合」同源。
///
private const string KpiEventFrom = """
FROM `mdp_std_action_event`
WHERE `tenant_id` = @TenantId
AND `biz_type` = 'crm_seorder'
AND `biz_id` IS NOT NULL
AND `action_code` IN ('S1_ORDER_REVIEW', 'S1_DELIVERY_CONFIRM')
""";
// ══════════════════════════════════════════════════════════════════════
// Stage-2 · 产品设计 KPI(dwd_product_design_kpi)
// ══════════════════════════════════════════════════════════════════════
///
/// Stage-2「产品设计」KPI。Grain = tenant_id + sales_order_id。
///
/// Start 直接取 dwd_order_review_kpi.review_end_time —— Stage-1 已有正式
/// DWD Authority,这里只消费,不重新实现 Stage-1 的配对逻辑。
///
/// End 取该订单 S2_DAILY_PLAN_RELEASE 的 FIRST SUCCESS
/// (ORDER BY action_end_time, action_run_id)的 action_end_time。
/// 将来同一订单再次下达也不会改写已成型的 Stage-2 End —— 首个里程碑一旦产生即恒定,
/// 故本 transform 是 deterministic / idempotent。
///
/// END_UNKNOWN 是本表存在的核心理由:1.0.532 之前没有释放事件,历史上已下达的订单
/// 只能证明「已下达」、无法证明「何时下达」。这类订单必须与「确实还没下达」的 IN_PROGRESS
/// 区分开,否则会把历史数据缺口误判成业务未完成。区分依据是中台里是否存在该订单
/// status='Y' 的日计划行 —— 只取这个布尔事实,绝不取它的时间。
///
/// MDP ONLY:全部取数来自 dwd_order_review_kpi / mdp_std_action_event /
/// mdp_std_operation_schedule / mdp_std_work_order_schedule / mdp_std_so,
/// 不 JOIN 任何源业务表。
///
/// 订单全集以 mdp_std_so 为准,而不是由释放事件驱动 ——
/// 否则「尚未开始评审」的订单会从档案里整个消失。
///
public async Task BuildProductDesignKpiAsync(
long tenantId, string batchId, DateTime now, CancellationToken cancellationToken = default)
{
if (tenantId <= 0) throw new InvalidOperationException("BuildProductDesignKpiAsync 需要有效 tenantId。");
cancellationToken.ThrowIfCancellationRequested();
var ps = new SugarParameter[]
{
new("@TenantId", tenantId),
new("@BatchId", batchId),
new("@Now", now)
};
// 逻辑计数:本轮参与 Stage-2 KPI transform 的 distinct 销售订单数(非 DB affected rows)
var logicalRows = await _db.Ado.GetIntAsync($"""
SELECT COUNT(*) FROM ({ProductDesignOrderScope}) o
""", ps);
const string sql = $"""
INSERT INTO `dwd_product_design_kpi`
(`tenant_id`, `sales_order_id`, `sales_order_no`, `review_end_time`, `release_action_run_id`,
`product_design_start_time`, `product_design_end_time`, `actual_hours`, `calculation_status`,
`release_success_count`, `release_failed_count`, `has_released_schedule`,
`calc_batch_id`, `calc_time`)
WITH orders AS ({ProductDesignOrderScope}),
-- 该订单的释放事件全集(用于计数;SUCCESS/FAILED 都留在事件层)
rel_all AS (
SELECT `tenant_id`, `biz_id`,
SUM(`action_status` = 'SUCCESS') AS succ_cnt,
SUM(`action_status` = 'FAILED') AS fail_cnt
FROM `mdp_std_action_event`
WHERE `tenant_id` = @TenantId
AND `action_code` = 'S2_DAILY_PLAN_RELEASE'
AND `biz_type` = 'crm_seorder'
AND `biz_id` IS NOT NULL
GROUP BY `tenant_id`, `biz_id`
),
-- FIRST SUCCESS:Stage-2 End 永远取首个成功释放,后续释放不改写历史
rel_rank AS (
SELECT `tenant_id`, `biz_id`, `action_run_id`, `action_end_time`,
ROW_NUMBER() OVER (PARTITION BY `tenant_id`, `biz_id`
ORDER BY `action_end_time` ASC, `action_run_id` ASC) AS rn
FROM `mdp_std_action_event`
WHERE `tenant_id` = @TenantId
AND `action_code` = 'S2_DAILY_PLAN_RELEASE'
AND `biz_type` = 'crm_seorder'
AND `biz_id` IS NOT NULL
AND `action_status` = 'SUCCESS'
AND `action_end_time` IS NOT NULL
),
rel1 AS (SELECT * FROM rel_rank WHERE rn = 1),
-- 历史「已下达」状态佐证:MDP-only ID 桥
-- operation_schedule → work_order_schedule.sales_order_entry_id → std_so.order_entry_id → order_id
-- 只产出布尔,不取任何时间(该层的 update_time 会被覆盖,不是下达时刻)
rel_sched AS (
SELECT DISTINCT so.`tenant_id`, so.`order_id`
FROM `mdp_std_operation_schedule` op
JOIN `mdp_std_work_order_schedule` wo
ON wo.`tenant_id` = op.`tenant_id` AND wo.`work_order` = op.`work_order`
JOIN `mdp_std_so` so
ON so.`tenant_id` = wo.`tenant_id` AND so.`order_entry_id` = wo.`sales_order_entry_id`
WHERE op.`tenant_id` = @TenantId
AND op.`status` = 'Y'
AND wo.`sales_order_entry_id` IS NOT NULL
AND so.`order_id` IS NOT NULL
)
SELECT
o.`tenant_id`,
o.`order_id`,
o.order_no,
k.`review_end_time`,
r.`action_run_id`,
k.`review_end_time` AS product_design_start_time,
r.`action_end_time` AS product_design_end_time,
CASE WHEN k.`review_end_time` IS NOT NULL
AND r.`action_end_time` IS NOT NULL
AND r.`action_end_time` > k.`review_end_time`
THEN ROUND(TIMESTAMPDIFF(SECOND, k.`review_end_time`, r.`action_end_time`) / 3600.0, 2)
END AS actual_hours,
CASE
-- 无 Start 无 End → 还没走到这一步
WHEN k.`review_end_time` IS NULL AND r.`action_end_time` IS NULL THEN 'NOT_STARTED'
-- 有 End 却没有 Start → 数据异常,不得倒推 Start
WHEN k.`review_end_time` IS NULL THEN 'INVALID_NO_START'
-- End <= Start → 时序异常,绝不 abs()/swap() 硬算正数
WHEN r.`action_end_time` IS NOT NULL
AND r.`action_end_time` <= k.`review_end_time` THEN 'INVALID_TIME_SEQUENCE'
WHEN r.`action_end_time` IS NOT NULL THEN 'OK'
-- 无释放事件但中台证明已下达 → 历史数据,时间 Authority 缺失
WHEN s.`order_id` IS NOT NULL THEN 'END_UNKNOWN'
ELSE 'IN_PROGRESS'
END AS calculation_status,
IFNULL(ra.succ_cnt, 0),
IFNULL(ra.fail_cnt, 0),
CASE WHEN s.`order_id` IS NOT NULL THEN 1 ELSE 0 END,
@BatchId, @Now
FROM orders o
LEFT JOIN `dwd_order_review_kpi` k
ON k.`tenant_id` = o.`tenant_id` AND k.`sales_order_id` = o.`order_id`
LEFT JOIN rel1 r ON r.`tenant_id` = o.`tenant_id` AND r.`biz_id` = o.`order_id`
LEFT JOIN rel_all ra ON ra.`tenant_id` = o.`tenant_id` AND ra.`biz_id` = o.`order_id`
LEFT JOIN rel_sched s ON s.`tenant_id` = o.`tenant_id` AND s.`order_id` = o.`order_id`
ON DUPLICATE KEY UPDATE
`sales_order_no` = VALUES(`sales_order_no`),
`review_end_time` = VALUES(`review_end_time`),
`release_action_run_id` = VALUES(`release_action_run_id`),
`product_design_start_time` = VALUES(`product_design_start_time`),
`product_design_end_time` = VALUES(`product_design_end_time`),
`actual_hours` = VALUES(`actual_hours`),
`calculation_status` = VALUES(`calculation_status`),
`release_success_count` = VALUES(`release_success_count`),
`release_failed_count` = VALUES(`release_failed_count`),
`has_released_schedule` = VALUES(`has_released_schedule`),
`calc_batch_id` = VALUES(`calc_batch_id`),
`calc_time` = VALUES(`calc_time`)
""";
// 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
await _db.Ado.ExecuteCommandAsync(sql, ps);
return logicalRows;
}
///
/// Stage-2 KPI 的订单全集:以 mdp_std_so 的订单头为准(一订单一行)。
/// 被计数查询与 CTE orders 共用。刻意不由释放事件驱动 ——
/// 那样「尚未开始评审」的订单会从档案里消失。
///
private const string ProductDesignOrderScope = """
SELECT `tenant_id`, `order_id`, MAX(`order_no`) AS order_no
FROM `mdp_std_so`
WHERE `tenant_id` = @TenantId
AND `order_id` IS NOT NULL
AND `order_id` <> 0
AND `deleted_flag` = 0
GROUP BY `tenant_id`, `order_id`
""";
}
///
/// 动作事件链路跑批结果。
///
/// 计数语义统一为「本轮参与处理的逻辑对象数」,不是数据库 affected rows。
/// 原因:INSERT … ON DUPLICATE KEY UPDATE 在 MySQL 下「插入计 1、更新计 2」,
/// 直接返回会让「什么都没变」的第二次跑批显示成处理量翻倍,属于虚报。
/// 因此每层都改为对该层的取数集合求 COUNT(DISTINCT 业务键):
/// 一个 action_run_id 最多计 1、一张销售订单最多计 1,重复 upsert 不翻倍。
///
public sealed class ActionEventSyncResult
{
public string BatchId { get; set; } = "";
public long TenantId { get; set; }
/// 本轮实际进入贴源处理链的 source event 数(增量拉取到的源行数)。
public int StageRows { get; set; }
/// 本轮被 RUNNING 回捞补偿重拉的 distinct action_run_id 数。
public int StageRepulledRunning { get; set; }
/// 本轮参与 STD transform 的 distinct action_run_id 数。
public int StandardRows { get; set; }
/// 本轮参与 KPI transform 的 distinct 销售订单数。
public int KpiRows { get; set; }
/// 本轮参与 Stage-2 产品设计 KPI transform 的 distinct 销售订单数。
public int ProductDesignKpiRows { get; set; }
}