merge: W-REAL R1(telemetry MQ 异步消费·灰度开关默认关)+ R6(D1 次留旁路聚合 V27)+ R2(积分充值 mock-gated SPI 骨架 V28)

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
lili 2026-07-02 20:51:44 -07:00
commit 4e1eafd634
37 changed files with 2628 additions and 80 deletions

View File

@ -100,6 +100,13 @@
<artifactId>huijing-spring-boot-starter-rpc</artifactId>
</dependency>
<!-- RocketMQ切片一 阶段〇 R1telemetry 批量事件入口改 MQ 异步消费RocketMQTemplate / @RocketMQMessageListener。
starter-mq 以 optional=true 声明本 starter不传递故本模块须显式直依赖版本由 huijing 依赖管理统一控制。 -->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>
<!-- 测试 -->
<dependency>
<groupId>com.wanxiang</groupId>

View File

@ -0,0 +1,40 @@
package com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat;
import com.wanxiang.huijing.framework.tenant.core.db.TenantBaseDO;
import com.baomidou.mybatisplus.annotation.KeySequence;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import lombok.EqualsAndHashCode;
import java.time.LocalDate;
/**
* 玩家首玩追踪 DO对应表 game_telemetry_player_first_play游戏 × 玩家 粒度
*
* <p>质量模型设计 §3.4 L3 留存结构落地件每玩家对每游戏一行记首玩自然日 + 是否在首玩日+1 自然日回访开局
* 继承 {@link TenantBaseDO} 自动携带 creator/create_time/updater/update_time/deleted + tenant_id
* 唯一约束 (game_id, player_key) 保证每玩家每游戏一行既是首玩落点也是 D1 去重幂等键
*
* <p>玩家标识 {@link #playerKey} 命名空间隔离登录 = {@code u:{userId}}未登录 = {@code a:{anonId}}
* 对齐 #5 契约 anonId 口径 userId anonId 字面碰撞
*
* @author 绘境AI
*/
@TableName("game_telemetry_player_first_play")
@KeySequence("game_telemetry_player_first_play_seq")
@Data
@EqualsAndHashCode(callSuper = true)
public class PlayerFirstPlayDO extends TenantBaseDO {
/** 记录 ID */
private Long id;
/** 游戏 ID= game_project.id */
private Long gameId;
/** 稳定玩家标识:登录 = u:{userId}、未登录 = a:{anonId}(命名空间隔离防碰撞) */
private String playerKey;
/** 该玩家对该游戏的首玩自然日(平台时区,与 game_telemetry_game_stat.stat_date 同日界) */
private LocalDate firstPlayDate;
/** 次留标记0 未回 / 1 已在首玩日+1 自然日回访开局CAS 置位一次、去重) */
private Integer d1Returned;
}

View File

@ -0,0 +1,37 @@
package com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat;
import com.wanxiang.huijing.framework.tenant.core.db.TenantBaseDO;
import com.baomidou.mybatisplus.annotation.KeySequence;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import lombok.EqualsAndHashCode;
import java.time.LocalDate;
/**
* 游戏维度 D1 次留聚合 DO对应表 game_telemetry_retention_stat游戏 × cohort日 粒度
*
* <p>质量模型设计 §3.4 L3 留存结构聚合件 (game_id, stat_date=cohort日) 聚合当日首玩新玩家数 new_player_count
* 与其中次日回访开局的 d1_retained_countD1 留存率 = d1_retained_count / new_player_count读时计算不落列
* 独立于 game_telemetry_game_stat不动既有聚合表结构quality_score 链路零风险
*
* @author 绘境AI
*/
@TableName("game_telemetry_retention_stat")
@KeySequence("game_telemetry_retention_stat_seq")
@Data
@EqualsAndHashCode(callSuper = true)
public class RetentionStatDO extends TenantBaseDO {
/** 聚合记录 ID */
private Long id;
/** 游戏 ID= game_project.id */
private Long gameId;
/** cohort 日 = 首玩自然日(平台时区) */
private LocalDate statDate;
/** 当日首玩新玩家数cohort 规模,玩家×游戏×日去重) */
private Long newPlayerCount;
/** 当日 cohort 中在次日回访开局的玩家数D1 留存分子) */
private Long d1RetainedCount;
}

View File

@ -0,0 +1,68 @@
package com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat;
import com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat.PlayerFirstPlayDO;
import com.wanxiang.huijing.framework.mybatis.core.mapper.BaseMapperX;
import com.wanxiang.huijing.framework.mybatis.core.query.LambdaQueryWrapperX;
import org.apache.ibatis.annotations.Insert;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Update;
import java.time.LocalDate;
/**
* 玩家首玩追踪 MapperD1 次留旁路游戏 × 玩家 粒度
*
* <p>并发安全靠两条原子语句{@link #insertIgnoreFirstPlay}INSERT IGNORE首玩幂等落点命中 uk_game_player 即忽略
* + {@link #markD1ReturnedIfUnset}CAS 置位d1_returned 01 恰一次审计/租户缺省列依赖 DDL 默认值MVP 单租户=0
* {@code GameStatMapper.insertOrAccumulate} 同范式
*
* @author 绘境AI
*/
@Mapper
public interface PlayerFirstPlayMapper extends BaseMapperX<PlayerFirstPlayDO> {
/**
* 首玩幂等落点INSERT IGNORE新玩家插入一行d1_returned=0命中 uk_game_player(game_id, player_key) 即忽略
*
* <p>返回受影响行数区分真首玩 vs 已存在1=真插入该玩家对该游戏首次开局 调用方据此累加 new_player_count
* 0=已存在同玩家此前已首玩 调用方转 D1 回访判定并发双首玩由 MySQL 在唯一键上原子裁决一方插入一方忽略
*
* @param gameId 游戏 ID
* @param playerKey 稳定玩家标识u:{userId} / a:{anonId}
* @param firstPlayDate 首玩自然日平台时区
* @return 受影响行数1=真首玩0=已存在含并发被忽略
*/
@Insert("INSERT IGNORE INTO game_telemetry_player_first_play "
+ "(game_id, player_key, first_play_date, d1_returned) "
+ "VALUES (#{gameId}, #{playerKey}, #{firstPlayDate}, 0)")
int insertIgnoreFirstPlay(@Param("gameId") Long gameId, @Param("playerKey") String playerKey,
@Param("firstPlayDate") LocalDate firstPlayDate);
/**
* 取玩家对某游戏的首玩行 D1 回访窗口用 first_play_date d1_returned
*
* @param gameId 游戏 ID
* @param playerKey 稳定玩家标识
* @return 首玩行不存在返回 null
*/
default PlayerFirstPlayDO selectByGameAndPlayer(Long gameId, String playerKey) {
return selectOne(new LambdaQueryWrapperX<PlayerFirstPlayDO>()
.eq(PlayerFirstPlayDO::getGameId, gameId)
.eq(PlayerFirstPlayDO::getPlayerKey, playerKey));
}
/**
* CAS 置位 D1 回访标记d1_returned 01恰一次去重仅当当前为 0 才置 1
*
* <p>返回受影响行数保证同一玩家次日多次开局只计一次 D11=本次首次置位调用方据此累加 d1_retained_count
* 0=已被置位重复回访去重跳过
*
* @param id 首玩行主键
* @return 受影响行数1=首次置位0=已置位去重
*/
@Update("UPDATE game_telemetry_player_first_play SET d1_returned = 1 "
+ "WHERE id = #{id} AND d1_returned = 0")
int markD1ReturnedIfUnset(@Param("id") Long id);
}

View File

@ -0,0 +1,67 @@
package com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat;
import com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat.RetentionStatDO;
import com.wanxiang.huijing.framework.mybatis.core.mapper.BaseMapperX;
import com.wanxiang.huijing.framework.mybatis.core.query.LambdaQueryWrapperX;
import org.apache.ibatis.annotations.Insert;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import java.time.LocalDate;
/**
* 游戏维度 D1 次留聚合 Mapper游戏 × cohort日 粒度
*
* <p>两条行级原子累加语句INSERTON DUPLICATE KEY UPDATE命中 uk_game_date 累加消除读--写并发丢增量
* {@code GameStatMapper.insertOrAccumulate} 同范式审计/租户缺省列依赖 DDL 默认值MVP 单租户=0
* 幂等边界cohort 计数的恰一次 first_play 表的 INSERT IGNORE / CAS 置位保证 RetentionAggregateService
* 本表只做纯加法
*
* @author 绘境AI
*/
@Mapper
public interface RetentionStatMapper extends BaseMapperX<RetentionStatDO> {
/**
* 累加当日首玩新玩家数cohort 规模 +1命中 (game_id, stat_date) new_player_count+1否则建行初值 (1,0)
*
* @param gameId 游戏 ID
* @param statDate cohort = 首玩自然日
* @return 受影响行数插入=1 / 累加=2仅日志排障用
*/
@Insert("INSERT INTO game_telemetry_retention_stat "
+ "(game_id, stat_date, new_player_count, d1_retained_count) "
+ "VALUES (#{gameId}, #{statDate}, 1, 0) "
+ "ON DUPLICATE KEY UPDATE new_player_count = new_player_count + 1")
int insertOrAccumulateNewPlayer(@Param("gameId") Long gameId, @Param("statDate") LocalDate statDate);
/**
* 累加当日 cohort D1 回访数分子 +1命中 (game_id, stat_date) d1_retained_count+1否则建行初值 (0,1)
*
* <p>常态下该 (game_id, cohort日) 行在首玩日已由 {@link #insertOrAccumulateNewPlayer} 建出故走累加分支
* VALUES(0,1) 插入分支为防御性兜底cohort 行意外缺失时不丢分子
*
* @param gameId 游戏 ID
* @param statDate cohort = 该玩家的首玩自然日
* @return 受影响行数插入=1 / 累加=2仅日志排障用
*/
@Insert("INSERT INTO game_telemetry_retention_stat "
+ "(game_id, stat_date, new_player_count, d1_retained_count) "
+ "VALUES (#{gameId}, #{statDate}, 0, 1) "
+ "ON DUPLICATE KEY UPDATE d1_retained_count = d1_retained_count + 1")
int insertOrAccumulateD1Retained(@Param("gameId") Long gameId, @Param("statDate") LocalDate statDate);
/**
* 游戏×cohort日 取次留聚合行对账/看板/回流读口
*
* @param gameId 游戏 ID
* @param statDate cohort
* @return 聚合行不存在返回 null
*/
default RetentionStatDO selectByGameAndDate(Long gameId, LocalDate statDate) {
return selectOne(new LambdaQueryWrapperX<RetentionStatDO>()
.eq(RetentionStatDO::getGameId, gameId)
.eq(RetentionStatDO::getStatDate, statDate));
}
}

View File

@ -1,51 +1,80 @@
package com.wanxiang.huijing.game.module.telemetry.mq.consumer;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.game.module.telemetry.mq.message.TelemetryEventMessage;
import com.wanxiang.huijing.game.module.telemetry.mq.producer.TelemetryEventProducer;
import com.wanxiang.huijing.game.module.telemetry.service.event.EventIngestService;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component;
import jakarta.annotation.Resource;
/**
* 遥测事件 MQ 消费者脊柱对接点骨架落原始事件 + 增量聚合 + 回灌
* 遥测事件 MQ 消费者切片一 阶段 R1异步消费事件信封落原始事件 + 增量聚合 + quality_score 回灌 feed
*
* 本类为外部依赖对接点骨架阶段不实现具体消费逻辑 mock不实现仅声明对接契约与 TODO保证集成可追溯
* 契约见 contracts/api-schemas/telemetry.yaml#x-mq-contracts.consumer
* - group=telemetry-event-consumertopic=telemetry-event无序并发消费聚合按 gameId 行级原子累加避免覆盖
* - 幂等键 = (traceId, event, ts)命中 game_telemetry_event.uk_dedup 即跳过INSERT IGNORE 语义
* 保证至少一次投递下不重复聚合
* - 默认重试 16 次仍失败转 DLQtelemetry-event-dlq T-TEL-19 数据质量监控告警
* <p><b>脊柱对接点落地</b>契约 {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.consumer}
* topic={@link TelemetryEventProducer#TOPIC}consumerGroup=telemetry-event-consumer并发消费无序
* 消费副作用逐字复用同步写核心 {@code EventIngestServiceImpl#ingestOne} {@link EventIngestService#consumeFromMq}
* 单条事务包裹<b>聚合行为与 quality_score 回灌 feed 行为与改造前逐字一致</b>只是执行位置从请求线程搬到消费线程
*
* 消费副作用待实现
* 1. game_telemetry_event原始事件幂等 TelemetryEventMapper insert依赖 uk_dedup 去重
* 2. 增量聚合 game_telemetry_game_statplay/like/quality 等按 gameId+statDate 行级累加唯一约束 uk_game_date 保证幂等
* 3. quality_score 计算运营质量分 0-100输入avgDurationMs/completedCount/reportCount/loadFailCount
* TODO 对接点计算公式待产品/数据侧确定后落地不在骨架阶段臆造公式
* 4. 回灌play_count/like_count projectgame_projectquality_score feed 排序信号
* TODO 对接点 project/feed -apiFeign就绪后在事务提交后异步回灌远程调用不放事务内
* <p><b>幂等</b>以信封 {@code eventId}uk_event_id为幂等键同步写核心先 selectByEventId 预查命中即跳过累加
* 叠加 uk_event_id 唯一约束硬后盾至少一次投递下重复消息不重复聚合契约 V10 uk_dedup 已被 uk_event_id 取代
*
* <p><b>并发上限{@code consumeThreadMax} 死参数坑</b>rocketmq-client 无界消费队列致 {@code consumeThreadMax}
* 不单独生效线程数永不超 core故与 {@code consumeThreadNumber} <b>同设为 8min=max</b>让消费并发上限真为 8
* MQ cap 投递消费速率不是在跑生成并发那种业务并发8 取值受 DB 连接池默认 10约束留头寸
* 每条消费=一次快 DB 写事务可按吞吐与池大小调改这里需同调池
*
* <p><b>重试与 DLQ外部交互纪律</b>消费抛异常则 RocketMQ 重投默认 16 仍失败转 DLQ%RETRY%/%DLQ%
* T-TEL-19 数据质量监控告警故对<b>瞬时故障</b>DB 抖动本方法<b>外抛让其重试</b>单条事务已回滚幂等保证重试不双计
* <b>毒消息</b>信封为空等结构性坏数据重投无益<b>吞掉正常返回</b>ack只告警留痕防无限重投风暴
*
* <p><b>装配条件</b>{@code telemetry.mq.enabled=true} 才装配与生产者同一灰度开关默认关时不起消费者不连 NameServer
*
* @author 绘境AI
*/
@Slf4j
@Component
public class TelemetryEventConsumer {
@ConditionalOnProperty(prefix = "telemetry.mq", name = "enabled", havingValue = "true")
@RocketMQMessageListener(
topic = TelemetryEventProducer.TOPIC,
consumerGroup = "telemetry-event-consumer",
consumeThreadNumber = 8, // min=max=8消费并发上限真为 8 DB 池默认 10 约束留头寸见类注释死参数坑
consumeThreadMax = 8 // rocketmq-client 无界队列致本参不单独生效须与 consumeThreadNumber 同值
)
public class TelemetryEventConsumer implements RocketMQListener<TelemetryEventMessage> {
/** 事件摄取 Service消费副作用逐字复用同步写核心保证聚合/回灌行为不变)。 */
@Resource
private EventIngestService eventIngestService;
/**
* TODO 对接点MQ 消费补全 RocketMQ 监听注解与消费方法
* huijing-spring-boot-starter-mq RocketMQ listener 就绪后按如下骨架实现
* 消费一条事件信封单条事务落库 + 聚合 + quality_score 回灌幂等
*
* @RocketMQMessageListener(topic = "telemetry-event", consumerGroup = "telemetry-event-consumer")
* public class Xxx implements RocketMQListener<EnvelopeMessage> {
* public void onMessage(EnvelopeMessage msg) {
* // 1) 幂等落原始表uk_dedup 去重INSERT IGNORE
* // 2) 处理状态机 0待聚合 1已聚合EventProcessStatusEnum 校验合法流转
* // 3) 增量聚合 game_telemetry_game_stat行级原子累加
* // 4) 计算 quality_score公式 TODO
* // 5) 事务提交后异步回灌 feed/projectFeign 调用不放事务内
* }
* }
*
* 骨架阶段仅记 INFO 占位说明对接点存在避免遗漏
* @param message MQ 消息体信封 + 上报时点 userId {@link TelemetryEventMessage}
*/
public TelemetryEventConsumer() {
log.info("[TelemetryEventConsumer] 骨架对接点已就位MQ 消费/聚合/回灌 待对接(见类注释 TODO");
@Override
public void onMessage(TelemetryEventMessage message) {
// 毒消息守卫结构性坏数据消息/信封/eventId 为空重投无益 吞掉正常返回ack只告警防无限重投
if (message == null || message.getEnvelope() == null) {
log.error("[telemetry-mq] 消息或信封为空,丢弃不重投 message={}", message);
return;
}
EnvelopeReqVO envelope = message.getEnvelope();
try {
// 复用同步写核心单条 @Transactional落原始事件eventId 幂等 增量聚合 quality_score 回灌 feed
eventIngestService.consumeFromMq(envelope, message.getUserId());
log.info("[telemetry-mq] 消费事件完成 eventId={} event={} traceId={}",
envelope.getEventId(), envelope.getEvent(), envelope.getTraceId());
} catch (Exception e) {
// 瞬时故障DB 抖动等外抛让 RocketMQ 重投单条事务已回滚幂等保证重试不双计16 次仍失败转 DLQ 告警
// 错误路径必须留痕创始人铁律
log.error("[telemetry-mq] 消费事件失败外抛触发重投幂等保证不双计≤16 次失败转 DLQeventId={} event={} traceId={}",
envelope.getEventId(), envelope.getEvent(), envelope.getTraceId(), e);
throw e;
}
}
}

View File

@ -0,0 +1,43 @@
package com.wanxiang.huijing.game.module.telemetry.mq.message;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import lombok.Data;
import java.io.Serializable;
/**
* 遥测事件 MQ 消息体切片一 阶段 R1批量事件入口改 MQ 异步消费
*
* <p>契约对齐 {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.producer.messageBody}核心是
* 逐条投递的事件信封 {@link EnvelopeReqVO}此外携带上报时点解析出的登录 {@code userId}传输关切
* 非契约字段异步消费脱离了 Web 请求线程无法再从 SecurityContext 取登录态故须在生产时随信封带上
* 保证消费侧算 user_id 兜底与审计列填充与同步写路径逐字一致 {@code EventIngestServiceImpl#consumeFromMq}
*
* <p>幂等真身 = 信封 {@code eventId}消费侧命中 {@code uk_event_id} 去重MQ 消息 keys traceId 仅便于排障对账
* rocketmq-spring 默认 Jackson 转换器序列化/反序列化往返字段均为 POJO无循环引用
*
* @author 绘境AI
*/
@Data
public class TelemetryEventMessage implements Serializable {
private static final long serialVersionUID = 1L;
/** 单条事件信封(幂等真身 eventId 在其内;消费侧据此落库 + 聚合 + 回灌)。 */
private EnvelopeReqVO envelope;
/**
* 上报时点解析的登录用户 ID匿名为 null异步消费无 Web 请求线程取不到登录态故随消息携带
* 供消费侧 resolveUserId 兜底 + 匿名审计列兜底判定与同步写路径口径一致
*/
private Long userId;
/** 空构造Jackson 反序列化需要)。 */
public TelemetryEventMessage() {
}
public TelemetryEventMessage(EnvelopeReqVO envelope, Long userId) {
this.envelope = envelope;
this.userId = userId;
}
}

View File

@ -0,0 +1,96 @@
package com.wanxiang.huijing.game.module.telemetry.mq.producer;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.game.module.telemetry.mq.message.TelemetryEventMessage;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.client.producer.SendStatus;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.apache.rocketmq.spring.support.RocketMQHeaders;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
/**
* 遥测事件 MQ 生产者切片一 阶段 R1上报通道轻量校验后逐条投递事件信封到 {@code telemetry-event} 主题
* 落库 + 聚合 + quality_score 回灌从请求线程内同步写改为消费侧异步完成幂等
*
* <p><b>契约对齐</b> {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.producer}topic={@link #TOPIC}
* tag=event|perf{@code /events/batch} event{@code /perf/beacon} perfkeys=traceId逐条投递orderly=false
*
* <p><b>外部交互纪律创始人铁律</b>syncSend 带发送超时 {@link #SEND_TIMEOUT_MS}RocketMQ 生产者自带
* 发送失败重试retryTimesWhenSendFailed 默认 2作瞬时抖动兜底<b>幂等</b>由消费侧 {@code eventId}uk_event_id去重
* 保证至少一次投递不重复聚合发送最终失败<b>不外抛</b>返回 {@code false} 交由调用方按受理失败(rejected)处置
* 由前端 SDK 分批重传作最外层兜底遥测=分析数据可容忍偶发丢失且客户端重试绝不 500 阻断上报
*
* <p><b>装配条件</b>{@code telemetry.mq.enabled=true} 才装配与消费者同一灰度开关默认关时本 Bean 缺席
* {@code EventIngestServiceImpl} {@code ObjectProvider} 取空回落原同步写路径现行行为逐字不变
*
* @author 绘境AI
*/
@Slf4j
@Component
@ConditionalOnProperty(prefix = "telemetry.mq", name = "enabled", havingValue = "true")
public class TelemetryEventProducer {
/** 遥测事件主题(与消费者 {@code TelemetryEventConsumer} 同名订阅;契约 telemetry-event。 */
public static final String TOPIC = "telemetry-event";
/** 事件通道 tag/events/batch。 */
public static final String TAG_EVENT = "event";
/** 性能兜底通道 tag/perf/beacon。 */
public static final String TAG_PERF = "perf";
/** 发送超时(毫秒)——外部交互必设超时,避免上报请求被 MQ 抖动无限拖住。 */
private static final long SEND_TIMEOUT_MS = 3000L;
/** RocketMQ 模板rocketmq-spring 自动装配name-server 见 application*.yaml → mini-infra 100.64.0.8:9876。 */
private final RocketMQTemplate rocketMQTemplate;
public TelemetryEventProducer(RocketMQTemplate rocketMQTemplate) {
this.rocketMQTemplate = rocketMQTemplate;
}
/**
* 投递一条事件信封到 {@link #TOPIC}同步发 + 超时成功 SEND_OK true否则/异常返 false
*
* @param envelope 已通过轻量校验的单条信封
* @param userId 上报时点登录用户 ID匿名为 null随消息携带供消费侧口径一致
* @param tag 消息 tag{@link #TAG_EVENT} / {@link #TAG_PERF}
* @return 是否投递成功true=broker 已确认 SEND_OKfalse=发送失败/异常调用方按 rejected 处置
*/
public boolean publish(EnvelopeReqVO envelope, Long userId, String tag) {
// 目的地格式 topic:tagrocketmq-spring 约定
String destination = TOPIC + ":" + tag;
try {
// 消息体 = 信封 + userIdkeys traceId便于全链路排障对账非幂等键幂等真身是 eventId
TelemetryEventMessage payload = new TelemetryEventMessage(envelope, userId);
MessageBuilder<TelemetryEventMessage> builder = MessageBuilder.withPayload(payload);
if (StringUtils.hasText(envelope.getTraceId())) {
builder.setHeader(RocketMQHeaders.KEYS, envelope.getTraceId());
}
Message<TelemetryEventMessage> message = builder.build();
SendResult result = rocketMQTemplate.syncSend(destination, message, SEND_TIMEOUT_MS);
boolean ok = result != null && result.getSendStatus() == SendStatus.SEND_OK;
if (ok) {
log.info("[telemetry-mq] 投递事件信号 topic={} tag={} eventId={} traceId={} msgId={}",
TOPIC, tag, envelope.getEventId(), envelope.getTraceId(),
result.getMsgId());
} else {
// 发送非 SEND_OK如刷盘/主从超时不外抛 false 交调用方按 rejected 处置客户端重传兜底
log.warn("[telemetry-mq] 投递事件非 SEND_OK topic={} tag={} eventId={} status={}",
TOPIC, tag, envelope.getEventId(), result == null ? null : result.getSendStatus());
}
return ok;
} catch (Exception e) {
// 外部交互错误路径必须留痕MQ 抖动/超时不外抛 500 阻断上报 false rejected前端 SDK 分批重传兜底
log.warn("[telemetry-mq] 投递事件失败(不外抛,按 rejected 处置由前端重传兜底topic={} tag={} eventId={} traceId={}",
TOPIC, tag, envelope.getEventId(), envelope.getTraceId(), e);
return false;
}
}
}

View File

@ -32,4 +32,16 @@ public interface EventIngestService {
*/
Boolean ingestBeacon(EnvelopeReqVO reqVO, Long userId);
/**
* MQ 消费侧单条摄取切片一 阶段 R1异步消费入口单条事务落库 + 聚合 + quality_score 回灌
*
* <p> {@code TelemetryEventConsumer} 调用消费线程MQ 关闭时也作 {@code ingestBatch}/{@code ingestBeacon}
* 的同步回落逐条入口经代理保证单条事务逐字复用同步写核心保证聚合/回灌行为与改造前一致
* 幂等由信封 eventIduk_event_id保证重复投递不重复聚合
*
* @param reqVO 单条信封已在生产侧通过轻量校验
* @param userId 上报时点登录用户 ID匿名为 null MQ 消息携带保证消费侧口径与同步写一致
*/
void consumeFromMq(EnvelopeReqVO reqVO, Long userId);
}

View File

@ -9,6 +9,7 @@ import com.wanxiang.huijing.game.module.telemetry.dal.mysql.event.TelemetryEvent
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.GameStatMapper;
import com.wanxiang.huijing.game.module.telemetry.enums.EventProcessStatusEnum;
import com.wanxiang.huijing.game.module.telemetry.enums.TelemetryEventEnum;
import com.wanxiang.huijing.game.module.telemetry.mq.producer.TelemetryEventProducer;
import com.wanxiang.huijing.game.module.telemetry.service.quality.AnomalyFilterService;
import com.wanxiang.huijing.game.module.feed.api.FeedApi;
import com.wanxiang.huijing.game.module.feed.dto.FeedRankUpsertReqDTO;
@ -18,6 +19,8 @@ import com.wanxiang.huijing.framework.common.util.servlet.ServletUtils;
import com.wanxiang.huijing.framework.security.core.LoginUser;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.context.annotation.Lazy;
import org.springframework.security.authentication.UsernamePasswordAuthenticationToken;
import org.springframework.security.core.context.SecurityContextHolder;
import org.springframework.stereotype.Service;
@ -99,61 +102,61 @@ public class EventIngestServiceImpl implements EventIngestService {
@Resource
private AnomalyFilterService anomalyFilterService;
/*
* TODO 对接点MQ 生产§3.5 仅留契约+TODO 不实现"上报快速 ACK + 异步消费聚合"
* 注入 RocketMQ 生产者投递 topic=telemetry-event幂等键=eventId对应 uk_event_id消费侧复用本类同步写聚合逻辑
* contracts/api-schemas/telemetry.yaml#x-mq-contractsMVP 同步写已满足闭环不引入 MQ
/**
* 遥测事件 MQ 生产者软注入切片一 阶段 R1{@code telemetry.mq.enabled=true} 才在席上报入口把
* 落库 + 聚合 + quality_score 回灌改为快速 ACK + 投递 {@code telemetry-event} 主题消费侧异步完成幂等
* <p>{@code enabled=false}默认时缺席{@link ObjectProvider#getIfAvailable()} 取空 回落原同步写路径
* 现行行为逐字不变既有单测据此仍绿 Mockito 单测不注入本 Provider{@link #resolveProducer()} 空值守卫回落同步写
*/
@Resource
private ObjectProvider<TelemetryEventProducer> telemetryEventProducerProvider;
/**
* 自身代理引用切片一 阶段 R1MQ 关闭回落同步写时逐条经代理调 {@link #consumeFromMq}保证 {@code @Transactional}
* 单条事务真生效Spring 自调用不走代理事务失效故显式经代理{@code @Lazy} 破本 Bean 自引用初始化环
* <p> Mockito 单测无 Spring 代理本字段为 null{@link #invokeConsumeFromMq}/{@link #ingestBatchSyncFallback}
* 空值守卫退化为直调 {@code this}事务注解在单测本就无 Spring 语义不影响 mapper 交互断言
*/
@Lazy
@Resource
private EventIngestService self;
/**
* D1 次留旁路聚合软注入W-REAL R6消费/同步写单条摄取后附带记一次开局到次留旁路 game_play_start
* <p>best-effort + 独立旁路两表不动既有 game_telemetry_game_stat 结构quality_score 链路零风险
* Mockito 单测不注入本 Provider{@link #resolveRetentionService()} 空值守卫跳过既有单测行为不变
*/
@Resource
private ObjectProvider<com.wanxiang.huijing.game.module.telemetry.service.retention.RetentionAggregateService> retentionAggregateServiceProvider;
// ============================== 上报入口MQ 异步路由 / 同步回落==============================
/**
* 批量摄取切片一 阶段 R1 路由MQ 开启则逐条轻量校验后投递 {@code telemetry-event}快速 ACK消费侧异步聚合
* MQ 关闭则回落原同步写逐条单事务经代理批量上限超限整体拒绝先于任何投递/落库与改造前口径一致
*/
@Override
@Transactional(rollbackFor = Exception.class) // §3.5 C5 同步写落库+聚合+回灌+置已聚合在同一本地事务任一步失败整批回滚
public EventBatchResultVO ingestBatch(EventBatchReqVO reqVO, Long userId) {
// 2026-06-10 e2e 二轮实证P0-1 补口匿名上报是多级联写事件 insert stat 原子累加 feed rank 回灌
// markAggregated UPDATE其中 markAggregated updateById 运行时实证仍把 updater=null 写进
// UPDATE SET M-b 三因叠加同形态逐写点 set 兜底易漏级联feed 模块回灌同线程同雷故在入口
// 对匿名请求注入系统身份 AigcGenerateExecutor.executeWithSystemIdentity 范本的 Web 线程版
// 使全链 DefaultDBFieldHandler 填充 creator/updater="0"finally 清理防 Web 线程上下文泄漏
boolean injected = injectSystemIdentityIfAnonymous(userId);
try {
// 批量上限超限整体拒绝统一错误码出口前端据此分批重传
if (reqVO.getBatch().size() > MAX_BATCH_SIZE) {
throw exception(TELEMETRY_BATCH_SIZE_EXCEEDED);
}
// 本次批量受理回执 traceId便于排障优先复用首条信封 traceId缺失则生成
String batchTraceId = resolveBatchTraceId(reqVO);
int accepted = 0;
int rejected = 0;
// 部分成功语义逐条轻量校验通过则同步写落库+聚合+回灌失败计 rejected不中断整批
for (EnvelopeReqVO envelope : reqVO.getBatch()) {
if (!isEnvelopeValid(envelope)) {
rejected++;
continue;
}
ingestOne(envelope, userId);
accepted++;
}
// 关键链路日志受理成功记 INFO含回执 traceId 与计数便于全链路排障
log.info("[ingestBatch] 批量同步写完成 batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected);
EventBatchResultVO result = new EventBatchResultVO();
result.setAccepted(accepted);
result.setRejected(rejected);
result.setTraceId(batchTraceId);
return result;
} finally {
if (injected) {
SecurityContextHolder.clearContext();
}
// 批量上限超限整体拒绝统一错误码出口前端据此分批重传先判任何投递/落库前
if (reqVO.getBatch().size() > MAX_BATCH_SIZE) {
throw exception(TELEMETRY_BATCH_SIZE_EXCEEDED);
}
TelemetryEventProducer producer = resolveProducer();
if (producer != null) {
// MQ 异步路径逐条校验 + 投递 event tag快速 ACK 返回 accepted/rejected不在请求线程做聚合
return dispatchBatchToMq(reqVO, userId, producer);
}
// 同步回落MQ 关闭 / 单测逐条单事务落库 + 聚合 + 回灌行为与改造前一致
return ingestBatchSyncFallback(reqVO, userId);
}
/**
* 兜底单条上报切片一 阶段 R1 路由校验失败恒返 true不阻塞前端卸载校验通过 MQ 开启则投递 perf tag
* 关闭则同步写fire-and-forget任何路径恒返 true
*/
@Override
@Transactional(rollbackFor = Exception.class) // 兜底通道单条同步写 ingestBatch 同口径落库+聚合+回灌原子
public Boolean ingestBeacon(EnvelopeReqVO reqVO, Long userId) {
// 兜底通道页面卸载场景校验失败不抛错返回受理成功不触发前端异常
// 校验通过则同步写若同步写抛错极少DB 故障由全局异常处理返回错误响应
// sendBeacon fire-and-forget 不读响应不阻塞卸载事务该回滚仍回滚保证不脏写
// 兜底通道页面卸载场景校验失败不抛错返回受理成功不触发前端异常
if (!isEnvelopeValid(reqVO)) {
log.warn("[ingestBeacon] 信封校验未通过,丢弃 eventId={}, event={}, traceId={}",
reqVO == null ? null : reqVO.getEventId(),
@ -161,11 +164,36 @@ public class EventIngestServiceImpl implements EventIngestService {
reqVO == null ? null : reqVO.getTraceId());
return Boolean.TRUE;
}
// 匿名注入系统身份 ingestBatch 同口径级联写审计列兜底finally 清理
TelemetryEventProducer producer = resolveProducer();
if (producer != null) {
// MQ 异步投递 perf tagfire-and-forget 不阻塞卸载投递失败由前端 sendBeacon 天然不重试性能埋点可容忍丢失
producer.publish(reqVO, userId, TelemetryEventProducer.TAG_PERF);
return Boolean.TRUE;
}
// 同步回落单条事务落库 + 聚合 + 回灌 ingestBatch 同口径
invokeConsumeFromMq(reqVO, userId);
return Boolean.TRUE;
}
/**
* MQ 消费侧单条摄取切片一 阶段 R1消费线程或同步回落逐条入口单条 {@code @Transactional} 落库 + 聚合 + 回灌
* 逐字复用同步写核心 {@link #ingestOne}聚合/回灌行为与改造前一致匿名注入系统身份补审计列finally 清理
* <p>2026-06-10 e2e P0-1 补口沿用匿名上报多级联写事件 insert / stat 原子累加 / feed rank 回灌 / markAggregated UPDATE
* DefaultDBFieldHandler 审计列填充统一取 "0"消费线程无 Web 请求上下文ServletUtils.getClientIP() null
* 由异常剔除桩按仅 anonId 维度处理本就 MQ-ready
*/
@Override
@Transactional(rollbackFor = Exception.class) // 单条同步写落库+聚合+回灌+置已聚合在同一本地事务任一步失败该条回滚幂等保证重试不双计
public void consumeFromMq(EnvelopeReqVO reqVO, Long userId) {
boolean injected = injectSystemIdentityIfAnonymous(userId);
try {
ingestOne(reqVO, userId);
return Boolean.TRUE;
// R6D1 次留旁路聚合 game_play_startbest-effort 内部吞异常 quality_score 主链同事务但失败不连累它
com.wanxiang.huijing.game.module.telemetry.service.retention.RetentionAggregateService retention =
resolveRetentionService();
if (retention != null) {
retention.recordIfPlayStart(reqVO, userId);
}
} finally {
if (injected) {
SecurityContextHolder.clearContext();
@ -173,6 +201,76 @@ public class EventIngestServiceImpl implements EventIngestService {
}
}
// ---------- 路由内部实现 ----------
/** 解析 MQ 生产者(软注入 + 空值守卫Provider 缺席/未注入enabled=false 或单测)→ null → 回落同步写。 */
private TelemetryEventProducer resolveProducer() {
return telemetryEventProducerProvider == null ? null : telemetryEventProducerProvider.getIfAvailable();
}
/** 解析次留聚合服务(软注入 + 空值守卫Provider 未注入(单测)→ null → 跳过次留旁路(既有行为不变)。 */
private com.wanxiang.huijing.game.module.telemetry.service.retention.RetentionAggregateService resolveRetentionService() {
return retentionAggregateServiceProvider == null ? null : retentionAggregateServiceProvider.getIfAvailable();
}
/**
* MQ 异步路径批量投递逐条轻量校验通过则投递 event tag成功计 accepted投递失败计 rejected 交前端重传
* 校验失败计 rejected返回 accepted/rejected + 回执 traceId部分成功语义非整批失败
*/
private EventBatchResultVO dispatchBatchToMq(EventBatchReqVO reqVO, Long userId, TelemetryEventProducer producer) {
String batchTraceId = resolveBatchTraceId(reqVO);
int accepted = 0;
int rejected = 0;
for (EnvelopeReqVO envelope : reqVO.getBatch()) {
if (!isEnvelopeValid(envelope)) {
rejected++;
continue;
}
// 投递成功=已受理(accepted)投递失败=按受理失败(rejected)前端分批重传兜底幂等由 eventId 保证不双计
if (producer.publish(envelope, userId, TelemetryEventProducer.TAG_EVENT)) {
accepted++;
} else {
rejected++;
}
}
log.info("[ingestBatch] 批量投递 MQ 完成 batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected);
EventBatchResultVO result = new EventBatchResultVO();
result.setAccepted(accepted);
result.setRejected(rejected);
result.setTraceId(batchTraceId);
return result;
}
/**
* 同步回落批量摄取MQ 关闭 / 单测逐条轻量校验通过则经代理调 {@link #consumeFromMq}单条事务失败计 rejected
* 与改造前的差异仅事务粒度整批逐条更贴契约非整批失败语义坏事件不连累好事件聚合/回灌逻辑逐字不变
*/
private EventBatchResultVO ingestBatchSyncFallback(EventBatchReqVO reqVO, Long userId) {
String batchTraceId = resolveBatchTraceId(reqVO);
int accepted = 0;
int rejected = 0;
for (EnvelopeReqVO envelope : reqVO.getBatch()) {
if (!isEnvelopeValid(envelope)) {
rejected++;
continue;
}
invokeConsumeFromMq(envelope, userId);
accepted++;
}
log.info("[ingestBatch] 批量同步写完成MQ 关闭回落batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected);
EventBatchResultVO result = new EventBatchResultVO();
result.setAccepted(accepted);
result.setRejected(rejected);
result.setTraceId(batchTraceId);
return result;
}
/** 经代理调 {@link #consumeFromMq}(保证单条 {@code @Transactional} 真生效);单测无代理时空值守卫直调 this。 */
private void invokeConsumeFromMq(EnvelopeReqVO envelope, Long userId) {
EventIngestService target = (self != null) ? self : this;
target.consumeFromMq(envelope, userId);
}
/**
* 匿名上报注入系统身份2026-06-10 e2e 二轮 P0-1 补口范本=AigcGenerateExecutor.executeWithSystemIdentity
*

View File

@ -0,0 +1,28 @@
package com.wanxiang.huijing.game.module.telemetry.service.retention;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
/**
* D1 次留旁路聚合 Service质量模型设计 §3.4 L3 留存结构
*
* <p>只吃 {@code game_play_start}开局事件追踪每玩家对每游戏的首玩自然日并在首玩日+1 自然日回访开局时
* 置位 D1 回访产出 (game_id, cohort日) 粒度的 new_player_count / d1_retained_count独立旁路两表
* 不动既有 game_telemetry_game_stat 结构与 quality_score 链路
*
* @author 绘境AI
*/
public interface RetentionAggregateService {
/**
* 若为 game_play_start 事件则记一次开局到次留旁路否则 no-op
*
* <p><b>零风险纪律</b>内部 best-effort任何解析/落库异常都在内部吞掉并告警绝不外抛
* 故与它同事务的 quality_score 主链event 落库 + game_stat 聚合 + feed 回灌不受次留失败连累
* 幂等首玩 INSERT IGNORED1 CAS 置位同玩家同日多次开局不重复计数
*
* @param envelope 单条信封已通过轻量校验
* @param userId 上报时点登录用户 ID匿名为 null与信封 user.userId 一起决定稳定玩家标识
*/
void recordIfPlayStart(EnvelopeReqVO envelope, Long userId);
}

View File

@ -0,0 +1,121 @@
package com.wanxiang.huijing.game.module.telemetry.service.retention;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat.PlayerFirstPlayDO;
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.PlayerFirstPlayMapper;
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.RetentionStatMapper;
import com.wanxiang.huijing.game.module.telemetry.enums.TelemetryEventEnum;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import java.time.Instant;
import java.time.LocalDate;
import java.time.ZoneId;
/**
* D1 次留旁路聚合实现质量模型设计 §3.4 L3
*
* <p>口径同一 game_id同一稳定玩家标识登录取 userId未登录取 anonId在首玩自然日平台时区
* game_telemetry_game_stat.stat_date 同日界之后第 1 个自然日内 1 game_play_start玩家×游戏×去重
* 增量落法首玩 INSERT IGNORE first_play + new_player_count+1次日回访 CAS 置位 d1_returned + d1_retained_count+1
*
* <p> quality_score 主链同事务但 best-effort所有异常内部吞掉告警不外抛次留失败不连累 event 落库/聚合/回灌
*
* @author 绘境AI
*/
@Slf4j
@Service
public class RetentionAggregateServiceImpl implements RetentionAggregateService {
@Resource
private PlayerFirstPlayMapper playerFirstPlayMapper;
@Resource
private RetentionStatMapper retentionStatMapper;
@Override
public void recordIfPlayStart(EnvelopeReqVO envelope, Long userId) {
// 只吃开局事件game_play_start其余事件 no-op
if (envelope == null || !TelemetryEventEnum.GAME_PLAY_START.getEvent().equals(envelope.getEvent())) {
return;
}
// 解析落点三要素gameId聚合落点playerKey稳定玩家标识playDate平台时区自然日
Long gameId = parseLong(envelope.getContext() == null ? null : envelope.getContext().getGameId());
String playerKey = resolvePlayerKey(envelope, userId);
if (gameId == null || playerKey == null || envelope.getTs() == null) {
// 无游戏归属 / 无稳定玩家标识 / 无时间戳 无法归 cohort跳过不报错
return;
}
LocalDate playDate = Instant.ofEpochMilli(envelope.getTs()).atZone(ZoneId.systemDefault()).toLocalDate();
try {
// 首玩幂等落点INSERT IGNORE 真插入(=1)即该玩家对该游戏首次开局 累加 cohort 新玩家数
int inserted = playerFirstPlayMapper.insertIgnoreFirstPlay(gameId, playerKey, playDate);
if (inserted > 0) {
retentionStatMapper.insertOrAccumulateNewPlayer(gameId, playDate);
log.info("[retention] 记首玩 gameId={} player={} firstPlayDate={}", gameId, playerKey, playDate);
return;
}
// 已有首玩记录判是否首玩日+1 自然日回访D1 窗口且尚未置位 CAS 置位 + 累加分子
PlayerFirstPlayDO row = playerFirstPlayMapper.selectByGameAndPlayer(gameId, playerKey);
if (row == null || row.getFirstPlayDate() == null) {
return; // 极端并发下取不到罕见跳过下次回访再判
}
boolean unset = row.getD1Returned() == null || row.getD1Returned() == 0;
boolean isD1Window = playDate.equals(row.getFirstPlayDate().plusDays(1));
if (unset && isD1Window) {
int marked = playerFirstPlayMapper.markD1ReturnedIfUnset(row.getId());
if (marked > 0) {
// cohort = 该玩家的首玩日分子累加落在首玩日的 cohort 行上
retentionStatMapper.insertOrAccumulateD1Retained(gameId, row.getFirstPlayDate());
log.info("[retention] 记 D1 回访 gameId={} player={} cohortDate={} returnDate={}",
gameId, playerKey, row.getFirstPlayDate(), playDate);
}
}
// 同日重复开局 / 更晚回访( D1 窗口) / 已置位 无操作去重
} catch (Exception e) {
// 零风险纪律次留旁路任何失败都不连累 quality_score 主链同事务内已完成的 event/聚合/回灌照常提交
// 错误路径必须留痕创始人铁律
log.warn("[retention] 次留旁路聚合失败best-effort不影响 quality_score 主链gameId={} player={} eventId={}",
gameId, playerKey, envelope.getEventId(), e);
}
}
/**
* 稳定玩家标识登录 = {@code u:{userId}}未登录 = {@code a:{anonId}}命名空间隔离防碰撞
* userId 优先取信封 user.userId回退上报时点登录 userId均无则取信封 user.anonId
*
* @return 玩家标识 userId 也无 anonId 时返回 null无法归属跳过
*/
private String resolvePlayerKey(EnvelopeReqVO envelope, Long userId) {
Long uid = userId;
if (envelope.getUser() != null && StringUtils.hasText(envelope.getUser().getUserId())) {
Long parsed = parseLong(envelope.getUser().getUserId());
if (parsed != null) {
uid = parsed;
}
}
if (uid != null) {
return "u:" + uid;
}
if (envelope.getUser() != null && StringUtils.hasText(envelope.getUser().getAnonId())) {
return "a:" + envelope.getUser().getAnonId();
}
return null;
}
/** 字符串安全转 Long空/非法返回 null。 */
private static Long parseLong(String s) {
if (!StringUtils.hasText(s)) {
return null;
}
try {
return Long.parseLong(s.trim());
} catch (NumberFormatException ex) {
return null;
}
}
}

View File

@ -0,0 +1,140 @@
package com.wanxiang.huijing.game.module.telemetry.mq;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.game.module.telemetry.mq.consumer.TelemetryEventConsumer;
import com.wanxiang.huijing.game.module.telemetry.mq.message.TelemetryEventMessage;
import com.wanxiang.huijing.game.module.telemetry.mq.producer.TelemetryEventProducer;
import com.wanxiang.huijing.game.module.telemetry.service.event.EventIngestService;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.apache.rocketmq.spring.support.RocketMQMessageConverter;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;
import org.mockito.ArgumentCaptor;
import org.springframework.test.util.ReflectionTestUtils;
import java.nio.charset.StandardCharsets;
import java.util.UUID;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
/**
* R1MQ 消费聚合真跑集成证据 mini-infra RocketMQ{@code 100.64.0.8:9876} · Tailscale 直连绕系统代理
* 跑一次 produce broker consume 转调 {@code consumeFromMq} 全链路坐实生产者/消费者/消息序列化对真 broker 成立
*
* <p><b>覆盖边界</b> broker 的投递/消费/往返序列化 + 两个生产类{@link TelemetryEventProducer} /
* {@link TelemetryEventConsumer}在真 broker 上的行为聚合落库consumeFromMqingestOnegame_telemetry_game_stat
* 的逻辑由 {@code EventIngestServiceImplMqPathTest#testConsumeFromMq_reusesIngestOne_aggregatesAndBackfills} 单测坐实
* 逐字复用 ingestOne不重造 IT mock EventIngestService 捕获转调不连 MySQL整条
* HTTP POSTMQ真库聚合行 的单体 e2e 需后端窗口MySQL/Redis/Nacos 全栈与阶段 C1/C2/C3 集成验证同window处置
*
* <p><b>默认跳过</b>{@code @EnabledIfSystemProperty(telemetry.mq.e2e=1)}常规 {@code mvn test} 不连 broker不跑本 IT
* 出证据时显式 {@code mvn test -Dtelemetry.mq.e2e=1 -Dtest=TelemetryMqRoundTripIT}
*
* @author 绘境AI
*/
@EnabledIfSystemProperty(named = "telemetry.mq.e2e", matches = "1")
class TelemetryMqRoundTripIT {
/** mini-infra NameServerTailscale 直连;无系统代理 env直连不被 fake-ip 拦)。 */
private static final String NAMESRV = System.getProperty("telemetry.mq.namesrv", "100.64.0.8:9876");
@Test
void produce_consume_roundTrip_againstRealBroker() throws Exception {
String eventId = "e2e-" + UUID.randomUUID();
String traceId = "trace-" + UUID.randomUUID();
// 独立测试消费组earliest 起消费防错过刚发的消息
String group = "telemetry-event-consumer-e2e-" + System.currentTimeMillis();
BlockingQueue<byte[]> received = new ArrayBlockingQueue<>(64);
ObjectMapper om = new ObjectMapper();
// push consumer订阅 telemetry-event收到即把 body 塞队列
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(group);
consumer.setNamesrvAddr(NAMESRV);
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
consumer.subscribe(TelemetryEventProducer.TOPIC, TelemetryEventProducer.TAG_EVENT);
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, ctx) -> {
msgs.forEach(m -> received.offer(m.getBody()));
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();
// 新消费组首启需等首次 rebalance + 路由拉取 establish CONSUME_FROM_LAST_OFFSET 的尾偏移
// 在此之前投递会被跳到尾而漏收新组 offset 竞态 8s 让其 settle 后再投
Thread.sleep(8000);
// producer构造 RocketMQTemplaterocketmq-client DefaultMQProducer TelemetryEventProducer
org.apache.rocketmq.client.producer.DefaultMQProducer mqProducer =
new org.apache.rocketmq.client.producer.DefaultMQProducer("telemetry-event-producer-e2e");
mqProducer.setNamesrvAddr(NAMESRV);
mqProducer.start();
RocketMQTemplate template = new RocketMQTemplate();
template.setProducer(mqProducer);
template.setMessageConverter(new RocketMQMessageConverter().getMessageConverter());
TelemetryEventProducer producer = new TelemetryEventProducer(template);
try {
EnvelopeReqVO env = new EnvelopeReqVO();
env.setEventId(eventId);
env.setEvent("game_play_end");
env.setSchemaVersion("v1");
env.setTs(System.currentTimeMillis());
env.setTraceId(traceId);
EnvelopeReqVO.ContextVO ctx = new EnvelopeReqVO.ContextVO();
ctx.setGameId("1024");
ctx.setVersionId("2048");
env.setContext(ctx);
// 1) 真投递到真 broker true broker 已确认 SEND_OK
boolean sent = producer.publish(env, 99L, TelemetryEventProducer.TAG_EVENT);
assertTrue(sent, "真 broker 应确认 SEND_OK");
// 2) broker 回投 扫描出我这条 eventId 匹配容忍 topic 上的其它并发流量20s
TelemetryEventMessage roundTrip = null;
long deadline = System.currentTimeMillis() + 20000;
while (System.currentTimeMillis() < deadline) {
byte[] body = received.poll(2, TimeUnit.SECONDS);
if (body == null) {
continue;
}
TelemetryEventMessage m = om.readValue(new String(body, StandardCharsets.UTF_8), TelemetryEventMessage.class);
if (m.getEnvelope() != null && eventId.equals(m.getEnvelope().getEventId())) {
roundTrip = m;
break;
}
}
assertNotNull(roundTrip, "应从真 broker 消费到刚投递的消息(按 eventId 匹配)");
// 3) 往返序列化保真信封逐字段一致
assertNotNull(roundTrip.getEnvelope(), "消息信封应往返完整");
assertEquals(eventId, roundTrip.getEnvelope().getEventId(), "eventId 往返一致(幂等真身)");
assertEquals("game_play_end", roundTrip.getEnvelope().getEvent());
assertEquals("1024", roundTrip.getEnvelope().getContext().getGameId(), "gameId 往返一致(聚合落点)");
assertEquals(99L, roundTrip.getUserId(), "userId 随消息携带、往返一致");
// 4) 真消费者类处理真 broker 消息 转调 consumeFromMq聚合入口逐字复用 ingestOne
TelemetryEventConsumer realConsumer = new TelemetryEventConsumer();
EventIngestService ingestMock = mock(EventIngestService.class);
ReflectionTestUtils.setField(realConsumer, "eventIngestService", ingestMock);
realConsumer.onMessage(roundTrip);
ArgumentCaptor<EnvelopeReqVO> envCaptor = ArgumentCaptor.forClass(EnvelopeReqVO.class);
verify(ingestMock, timeout(2000)).consumeFromMq(envCaptor.capture(), eq(99L));
assertEquals(eventId, envCaptor.getValue().getEventId(), "消费者应把真 broker 信封转调聚合入口");
} finally {
consumer.shutdown();
mqProducer.shutdown();
}
}
}

View File

@ -0,0 +1,81 @@
package com.wanxiang.huijing.game.module.telemetry.mq.consumer;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.game.module.telemetry.mq.message.TelemetryEventMessage;
import com.wanxiang.huijing.game.module.telemetry.service.event.EventIngestService;
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
import org.junit.jupiter.api.Test;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
/**
* {@link TelemetryEventConsumer} 单元测试 Mockito不触真 MQ把守 R1 消费者转调 + 毒消息守卫 + 重试语义
*
* <p>覆盖 正常消费转调 {@code consumeFromMq(envelope,userId)} 毒消息消息/信封为空吞掉不转调不外抛ack 防无限重投
* 瞬时故障{@code consumeFromMq} 抛异常外抛让 RocketMQ 重投16 次转 DLQ幂等保证不双计
*
* @author 绘境AI
*/
class TelemetryEventConsumerTest extends BaseMockitoUnitTest {
@Mock
private EventIngestService eventIngestService;
@InjectMocks
private TelemetryEventConsumer consumer;
private static EnvelopeReqVO envelope() {
EnvelopeReqVO env = new EnvelopeReqVO();
env.setEventId("evt-1");
env.setEvent("game_play_end");
env.setSchemaVersion("v1");
env.setTs(1733600000000L);
env.setTraceId("trace-abc");
return env;
}
/** 正常路:消费转调 consumeFromMq(envelope, userId),逐字复用同步写核心。 */
@Test
void testOnMessage_dispatchesToConsumeFromMq() {
EnvelopeReqVO env = envelope();
TelemetryEventMessage msg = new TelemetryEventMessage(env, 99L);
assertDoesNotThrow(() -> consumer.onMessage(msg));
verify(eventIngestService).consumeFromMq(eq(env), eq(99L));
}
/** 毒消息守卫①:消息为 null → 吞掉、不转调、不外抛ack 防无限重投)。 */
@Test
void testOnMessage_nullMessage_swallowedNoDispatch() {
assertDoesNotThrow(() -> consumer.onMessage(null));
verify(eventIngestService, never()).consumeFromMq(any(), any());
}
/** 毒消息守卫②:信封为 null → 吞掉、不转调、不外抛。 */
@Test
void testOnMessage_nullEnvelope_swallowedNoDispatch() {
TelemetryEventMessage msg = new TelemetryEventMessage(null, 99L);
assertDoesNotThrow(() -> consumer.onMessage(msg));
verify(eventIngestService, never()).consumeFromMq(any(), any());
}
/** 瞬时故障consumeFromMq 抛异常 → onMessage 外抛(触发 RocketMQ 重投;幂等保证重试不双计)。 */
@Test
void testOnMessage_transientFailure_rethrowsForRetry() {
EnvelopeReqVO env = envelope();
TelemetryEventMessage msg = new TelemetryEventMessage(env, null);
doThrow(new RuntimeException("DB down")).when(eventIngestService).consumeFromMq(eq(env), isNull());
assertThrows(RuntimeException.class, () -> consumer.onMessage(msg));
}
}

View File

@ -0,0 +1,94 @@
package com.wanxiang.huijing.game.module.telemetry.mq.producer;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.client.producer.SendStatus;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.springframework.messaging.Message;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
/**
* {@link TelemetryEventProducer} 单元测试 Mockito不触真 MQ把守 R1 生产者投递契约与外部交互纪律
*
* <p>覆盖 正常投递到 {@code telemetry-event:event} true SEND_OK 判定 syncSend 抛异常时
* {@code publish} <b>不外抛</b> false交调用方按 rejected 处置前端重传兜底绝不 500 阻断上报
* SEND_OK false keys traceId
*
* @author 绘境AI
*/
class TelemetryEventProducerTest extends BaseMockitoUnitTest {
@Mock
private RocketMQTemplate rocketMQTemplate;
@InjectMocks
private TelemetryEventProducer producer;
/** 构造一条最简合法信封(含 eventId/traceId。 */
private static EnvelopeReqVO envelope() {
EnvelopeReqVO env = new EnvelopeReqVO();
env.setEventId("evt-1");
env.setEvent("game_play_start");
env.setSchemaVersion("v1");
env.setTs(1733600000000L);
env.setTraceId("trace-abc");
return env;
}
/** 正常路:投递到 topic:tag=telemetry-event:eventSEND_OK → 返 truekeys=traceId。 */
@Test
@SuppressWarnings("unchecked")
void testPublish_sendOk_returnsTrueAndKeysTraceId() {
SendResult ok = new SendResult();
ok.setSendStatus(SendStatus.SEND_OK);
ok.setMsgId("mid-1");
when(rocketMQTemplate.syncSend(eq(TelemetryEventProducer.TOPIC + ":" + TelemetryEventProducer.TAG_EVENT),
(Message<?>) org.mockito.ArgumentMatchers.any(), anyLong())).thenReturn(ok);
boolean sent = producer.publish(envelope(), 99L, TelemetryEventProducer.TAG_EVENT);
assertTrue(sent, "SEND_OK 应返回 true");
// 校验投递目的地 + keys header = traceId
ArgumentCaptor<Message<?>> msgCaptor = ArgumentCaptor.forClass(Message.class);
verify(rocketMQTemplate).syncSend(eq("telemetry-event:event"), msgCaptor.capture(), anyLong());
Object keys = msgCaptor.getValue().getHeaders().get("KEYS");
assertTrue("trace-abc".equals(keys), "MQ keys 应置 traceId 便于对账");
}
/** best-effort核心断言syncSend 抛异常 → publish 吞异常不外抛、返 false前端重传兜底不 500 阻断上报)。 */
@Test
void testPublish_bestEffort_swallowsExceptionReturnsFalse() {
doThrow(new RuntimeException("MQ down")).when(rocketMQTemplate)
.syncSend(org.mockito.ArgumentMatchers.anyString(),
(Message<?>) org.mockito.ArgumentMatchers.any(), anyLong());
boolean sent = assertDoesNotThrow(() -> producer.publish(envelope(), null, TelemetryEventProducer.TAG_EVENT));
assertFalse(sent, "发送异常应返回 false按 rejected 处置)");
}
/** 非 SEND_OK如刷盘超时→ 返 false不外抛。 */
@Test
@SuppressWarnings("unchecked")
void testPublish_nonSendOk_returnsFalse() {
SendResult flushTimeout = new SendResult();
flushTimeout.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
when(rocketMQTemplate.syncSend(org.mockito.ArgumentMatchers.anyString(),
(Message<?>) org.mockito.ArgumentMatchers.any(), anyLong())).thenReturn(flushTimeout);
boolean sent = producer.publish(envelope(), 1L, TelemetryEventProducer.TAG_PERF);
assertFalse(sent, "非 SEND_OK 应返回 false");
}
}

View File

@ -0,0 +1,188 @@
package com.wanxiang.huijing.game.module.telemetry.service.event;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EventBatchReqVO;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EventBatchResultVO;
import com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat.GameStatDO;
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.event.TelemetryEventMapper;
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.GameStatMapper;
import com.wanxiang.huijing.game.module.telemetry.mq.producer.TelemetryEventProducer;
import com.wanxiang.huijing.game.module.telemetry.service.quality.AnomalyFilterService;
import com.wanxiang.huijing.game.module.feed.api.FeedApi;
import com.wanxiang.huijing.game.module.feed.dto.FeedRankUpsertReqDTO;
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
import org.junit.jupiter.api.Test;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.springframework.beans.factory.ObjectProvider;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
/**
* {@link EventIngestServiceImpl} MQ 异步路径单测切片一 阶段 R1把守生产者在席=投递不同步写路由
* 消费入口 consumeFromMq 逐字复用 ingestOne 聚合/回灌 {@code EventIngestServiceImplTest}同步回落路径
* 既有断言逐字不变互补共同坐实聚合行为与 quality_score 回灌行为不变仅执行位置从请求线程搬到消费线程
*
* @author 绘境AI
*/
class EventIngestServiceImplMqPathTest extends BaseMockitoUnitTest {
@InjectMocks
private EventIngestServiceImpl eventIngestService;
@Mock
private TelemetryEventMapper telemetryEventMapper;
@Mock
private GameStatMapper gameStatMapper;
@Mock
private FeedApi feedApi;
@Mock
private AnomalyFilterService anomalyFilterService;
@Mock
private ObjectProvider<TelemetryEventProducer> telemetryEventProducerProvider; // MQ 生产者软注入
@Mock
private TelemetryEventProducer telemetryEventProducer;
/** MQ 开启ingestBatch 逐条投递、不在请求线程做同步写DB/feed 零交互)。 */
@Test
void testIngestBatch_mqEnabled_producesNotSync() {
when(telemetryEventProducerProvider.getIfAvailable()).thenReturn(telemetryEventProducer);
when(telemetryEventProducer.publish(any(), any(), eq(TelemetryEventProducer.TAG_EVENT))).thenReturn(true);
EventBatchReqVO reqVO = new EventBatchReqVO();
List<EnvelopeReqVO> batch = new ArrayList<>();
batch.add(validEnvelope("game_play_start", "t1")); // 合法 投递
batch.add(validEnvelope("like", "t2")); // 合法 投递
batch.add(validEnvelope("not_registered_evt", "t3")); // 未登记 rejected不投递
reqVO.setBatch(batch);
EventBatchResultVO result = eventIngestService.ingestBatch(reqVO, 99L);
assertEquals(2, result.getAccepted());
assertEquals(1, result.getRejected());
assertEquals("t1", result.getTraceId());
// 逐条投递 event tag 两次
verify(telemetryEventProducer, times(2)).publish(any(), eq(99L), eq(TelemetryEventProducer.TAG_EVENT));
// 关键请求线程内不做同步写落库/聚合/回灌全在消费侧
verifyNoInteractions(telemetryEventMapper);
verifyNoInteractions(gameStatMapper);
verifyNoInteractions(feedApi);
}
/** MQ 开启但投递失败:该条计 rejected前端重传兜底不误报 accepted。 */
@Test
void testIngestBatch_mqEnabled_publishFailCountedRejected() {
when(telemetryEventProducerProvider.getIfAvailable()).thenReturn(telemetryEventProducer);
when(telemetryEventProducer.publish(any(), any(), anyString())).thenReturn(false); // 投递恒失败
EventBatchReqVO reqVO = new EventBatchReqVO();
List<EnvelopeReqVO> batch = new ArrayList<>();
batch.add(validEnvelope("game_play_start", "t1"));
reqVO.setBatch(batch);
EventBatchResultVO result = eventIngestService.ingestBatch(reqVO, 99L);
assertEquals(0, result.getAccepted());
assertEquals(1, result.getRejected(), "投递失败应计 rejected 交前端重传");
}
/** MQ 开启beacon 投递 perf tag、恒返 true、不同步写。 */
@Test
void testIngestBeacon_mqEnabled_producesPerfTag() {
when(telemetryEventProducerProvider.getIfAvailable()).thenReturn(telemetryEventProducer);
when(telemetryEventProducer.publish(any(), any(), eq(TelemetryEventProducer.TAG_PERF))).thenReturn(true);
Boolean ok = eventIngestService.ingestBeacon(validEnvelope("perf_first_screen", "t-beacon"), 99L);
assertTrue(ok);
verify(telemetryEventProducer).publish(any(), eq(99L), eq(TelemetryEventProducer.TAG_PERF));
verifyNoInteractions(telemetryEventMapper);
}
/**
* 消费入口 consumeFromMq 逐字复用 ingestOnegame_play_end(completed) 原子累加 重读算分(60.00) 回灌 feed
* 坐实聚合/回灌行为不变 EventIngestServiceImplTest#testIngestOne_playEndTriggersQualityRefresh 同口径
* 只是入口从 ingestBatch 换成 consumeFromMq消费线程入口
*/
@Test
void testConsumeFromMq_reusesIngestOne_aggregatesAndBackfills() {
Map<String, Object> props = new HashMap<>();
props.put("completed", true);
props.put("duration_ms", 5000);
EnvelopeReqVO env = gameEnvelope("game_play_end", "t1", "1024", "2048", props);
when(telemetryEventMapper.selectByEventId(env.getEventId())).thenReturn(null);
GameStatDO accumulated = new GameStatDO();
accumulated.setId(7L);
accumulated.setPlayCount(0L);
accumulated.setPlayEndCount(1L);
accumulated.setCompletedCount(1L);
accumulated.setTotalDurationMs(5000L);
accumulated.setLikeCount(0L);
accumulated.setShareCount(0L);
when(gameStatMapper.selectByGameAndDate(eq(1024L), any(LocalDate.class))).thenReturn(accumulated);
eventIngestService.consumeFromMq(env, 99L);
// 原子累加play_end=1, completed=1, duration=5000
verify(gameStatMapper).insertOrAccumulate(eq(1024L), any(LocalDate.class),
eq(0L), eq(1L), eq(1L), eq(5000L), eq(0L), eq(0L));
// 算分写回 + 回灌 feedQUALITY_REFRESH60.0
verify(gameStatMapper).updateById(any(GameStatDO.class));
verify(feedApi).upsertRank(any(FeedRankUpsertReqDTO.class));
// 幂等预查发生重放同 eventId 不双计的第一道网
verify(telemetryEventMapper).selectByEventId(env.getEventId());
}
/** consumeFromMq 幂等:同 eventId 已落库 → 跳过累加/回灌(重复投递不双计)。 */
@Test
void testConsumeFromMq_idempotentSkipOnDuplicate() {
EnvelopeReqVO env = gameEnvelope("game_play_start", "t1", "1024", "2048", null);
when(telemetryEventMapper.selectByEventId(env.getEventId()))
.thenReturn(new com.wanxiang.huijing.game.module.telemetry.dal.dataobject.event.TelemetryEventDO());
eventIngestService.consumeFromMq(env, 99L);
verify(telemetryEventMapper, never()).insert(any(com.wanxiang.huijing.game.module.telemetry.dal.dataobject.event.TelemetryEventDO.class));
verifyNoInteractions(gameStatMapper);
verifyNoInteractions(feedApi);
}
// ============================== 测试夹具 ==============================
private static EnvelopeReqVO validEnvelope(String event, String traceId) {
EnvelopeReqVO envelope = new EnvelopeReqVO();
envelope.setEventId(UUID.randomUUID().toString());
envelope.setEvent(event);
envelope.setSchemaVersion("v1");
envelope.setTs(1733600000000L);
envelope.setTraceId(traceId);
return envelope;
}
private static EnvelopeReqVO gameEnvelope(String event, String traceId, String gameId, String versionId, Map<String, Object> props) {
EnvelopeReqVO envelope = validEnvelope(event, traceId);
EnvelopeReqVO.ContextVO ctx = new EnvelopeReqVO.ContextVO();
ctx.setGameId(gameId);
ctx.setVersionId(versionId);
envelope.setContext(ctx);
envelope.setProps(props);
return envelope;
}
}

View File

@ -0,0 +1,191 @@
package com.wanxiang.huijing.game.module.telemetry.service.retention;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.Statement;
import java.time.LocalDate;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* R6D1 次留聚合真跑集成证据 mini-infra MySQL{@code 100.64.0.8:3306} · Tailscale 直连JDBC 不走 HTTP 代理
* throwaway 库里跑<b> V27 迁移文件的 DDL</b> + 一个两日 D1 场景 new_player_count/d1_retained_count 计数正确
*
* <p><b>覆盖边界</b> MySQL V27 DDL 建表可跑双验 R2 单体 Flyway 前的 DDL 合法性 mapper 三条语句
* INSERT IGNORE 首玩 / ON DUP KEY 累加 / CAS 置位 D1行为正确 D1 cohort 计数首玩日+1 去重正确
* 服务分支逻辑首玩/D1窗口/去重 {@link RetentionAggregateServiceImplTest} 单测坐实 IT 用真 SQL 复演其决策序列
* 全程 throwaway {@code telemetry_r6_e2e} /隔离不碰真库真表不与 Flyway 冲突
*
* <p><b>默认跳过</b>{@code @EnabledIfSystemProperty(telemetry.retention.e2e=1)}常规 {@code mvn test} 不连 MySQL
* 出证据时 {@code mvn test -Dtelemetry.retention.e2e=1 -Dtest=RetentionAggregateRealDbIT}
*
* @author 绘境AI
*/
@EnabledIfSystemProperty(named = "telemetry.retention.e2e", matches = "1")
class RetentionAggregateRealDbIT {
private static final String HOST = System.getProperty("telemetry.retention.dbhost", "100.64.0.8");
private static final String USER = System.getProperty("telemetry.retention.dbuser", "root");
// 内网 MVP 凭据铁律授权入仓 docs/内网凭据与端点.md可经 -D 覆盖
private static final String PASS = System.getProperty("telemetry.retention.dbpass", "ZRH3jwYLOrntBcTAw29MW9BP");
private static final String DB = "telemetry_r6_e2e";
private static final long G = 999_999_001L; // 测试游戏 IDthrowaway 无关真数据
private static final LocalDate D = LocalDate.of(2026, 7, 2);
private String baseUrl() {
// JDBC/TCP 直连 Tailscale不经 HTTP 代理allowPublicKeyRetrieval + useSSL=false 内网 MVP
return "jdbc:mysql://" + HOST + ":3306/?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=UTC";
}
@Test
void v27Ddl_and_d1CohortCounting_onRealMysql() throws Exception {
// 内网直连绕系统 SOCKS 代理macOS SOCKSEnable127.0.0.1:7897MySQL Connector/J 默认走 JVM socks
// 直连 Tailscale MySQL 必须显式禁栽过的确切坑不禁则握手期 EOFread 0 bytes
System.setProperty("socksProxyHost", "");
System.setProperty("java.net.useSystemProxies", "false");
try (Connection conn = DriverManager.getConnection(baseUrl(), USER, PASS)) {
try (Statement st = conn.createStatement()) {
st.execute("DROP DATABASE IF EXISTS " + DB);
st.execute("CREATE DATABASE " + DB + " DEFAULT CHARACTER SET utf8mb4");
st.execute("USE " + DB);
}
// 跑真 V27 迁移文件的 DDL DDL 在真 MySQL 合法可建
runV27Ddl(conn);
// 两日 D1 场景复演 RetentionAggregateService 的决策序列 mapper 同款 SQL
// Day D玩家 u:1u:2 首玩 new_player_count=2
// Day D+1u:1 回访(D1 命中)u:2 不回 d1_retained_count=1
recordPlayStart(conn, "u:1", D); // 首玩
recordPlayStart(conn, "u:2", D); // 首玩
recordPlayStart(conn, "u:1", D.plusDays(1)); // D1 回访
// u:2 不回访 不产生 D+1 事件
// 断言 cohort 计数
try (PreparedStatement ps = conn.prepareStatement(
"SELECT new_player_count, d1_retained_count FROM game_telemetry_retention_stat "
+ "WHERE game_id=? AND stat_date=?")) {
ps.setLong(1, G);
ps.setObject(2, D);
try (ResultSet rs = ps.executeQuery()) {
assertTrue(rs.next(), "cohort 行应存在");
assertEquals(2L, rs.getLong(1), "首玩新玩家数应为 2u:1 + u:2");
assertEquals(1L, rs.getLong(2), "D1 回访数应为 1仅 u:1 次日回访)");
}
}
// 复核首玩表u:1 已置位 d1_returned=1u:2 未置位=0
assertEquals(1, queryD1Returned(conn, "u:1"), "u:1 应被置位 D1 回访");
assertEquals(0, queryD1Returned(conn, "u:2"), "u:2 未回访、d1_returned 应为 0");
} finally {
// 清理 throwaway 隔离不留痕
try (Connection conn = DriverManager.getConnection(baseUrl(), USER, PASS);
Statement st = conn.createStatement()) {
st.execute("DROP DATABASE IF EXISTS " + DB);
}
}
}
/** 读并执行真 V27 迁移文件的 DDL去 -- 注释行,按 ; 切分执行)。 */
private void runV27Ddl(Connection conn) throws Exception {
// 相对模块 basedirsurefire 工作目录定位 huijing-server 下的迁移文件
Path migration = Paths.get(System.getProperty("user.dir"), "..", "..",
"huijing-server", "src", "main", "resources", "db", "migration",
"V27.0.0__create_game_telemetry_retention.sql");
assertTrue(Files.exists(migration), "应找到 V27 迁移文件:" + migration.toAbsolutePath().normalize());
String sql = Files.readString(migration);
// 去掉整行 -- 注释后按 ; 切分
StringBuilder cleaned = new StringBuilder();
for (String line : sql.split("\n")) {
String trimmed = line.trim();
if (trimmed.startsWith("--")) {
continue;
}
cleaned.append(line).append('\n');
}
try (Statement st = conn.createStatement()) {
for (String stmt : cleaned.toString().split(";")) {
if (stmt.trim().isEmpty()) {
continue;
}
st.execute(stmt);
}
}
}
/** 复演 service.recordIfPlayStart 决策序列(真 SQL首玩累加 new_playerD1 窗口 CAS 置位 + 累加分子。 */
private void recordPlayStart(Connection conn, String playerKey, LocalDate playDate) throws Exception {
// INSERT IGNORE 首玩落点
int inserted;
try (PreparedStatement ps = conn.prepareStatement(
"INSERT IGNORE INTO game_telemetry_player_first_play "
+ "(game_id, player_key, first_play_date, d1_returned) VALUES (?,?,?,0)")) {
ps.setLong(1, G);
ps.setString(2, playerKey);
ps.setObject(3, playDate);
inserted = ps.executeUpdate();
}
if (inserted > 0) {
accumulate(conn, "INSERT INTO game_telemetry_retention_stat "
+ "(game_id, stat_date, new_player_count, d1_retained_count) VALUES (?,?,1,0) "
+ "ON DUPLICATE KEY UPDATE new_player_count = new_player_count + 1", playDate);
return;
}
// 已存在 D1 窗口 + CAS 置位
Long rowId = null;
LocalDate firstPlayDate = null;
int d1 = 0;
try (PreparedStatement ps = conn.prepareStatement(
"SELECT id, first_play_date, d1_returned FROM game_telemetry_player_first_play "
+ "WHERE game_id=? AND player_key=?")) {
ps.setLong(1, G);
ps.setString(2, playerKey);
try (ResultSet rs = ps.executeQuery()) {
if (rs.next()) {
rowId = rs.getLong(1);
firstPlayDate = rs.getObject(2, LocalDate.class);
d1 = rs.getInt(3);
}
}
}
if (rowId != null && d1 == 0 && firstPlayDate != null && playDate.equals(firstPlayDate.plusDays(1))) {
int marked;
try (PreparedStatement ps = conn.prepareStatement(
"UPDATE game_telemetry_player_first_play SET d1_returned=1 WHERE id=? AND d1_returned=0")) {
ps.setLong(1, rowId);
marked = ps.executeUpdate();
}
if (marked > 0) {
accumulate(conn, "INSERT INTO game_telemetry_retention_stat "
+ "(game_id, stat_date, new_player_count, d1_retained_count) VALUES (?,?,0,1) "
+ "ON DUPLICATE KEY UPDATE d1_retained_count = d1_retained_count + 1", firstPlayDate);
}
}
}
private void accumulate(Connection conn, String sql, LocalDate statDate) throws Exception {
try (PreparedStatement ps = conn.prepareStatement(sql)) {
ps.setLong(1, G);
ps.setObject(2, statDate);
ps.executeUpdate();
}
}
private int queryD1Returned(Connection conn, String playerKey) throws Exception {
try (PreparedStatement ps = conn.prepareStatement(
"SELECT d1_returned FROM game_telemetry_player_first_play WHERE game_id=? AND player_key=?")) {
ps.setLong(1, G);
ps.setString(2, playerKey);
try (ResultSet rs = ps.executeQuery()) {
return rs.next() ? rs.getInt(1) : -1;
}
}
}
}

View File

@ -0,0 +1,184 @@
package com.wanxiang.huijing.game.module.telemetry.service.retention;
import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
import com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat.PlayerFirstPlayDO;
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.PlayerFirstPlayMapper;
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.RetentionStatMapper;
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
import org.junit.jupiter.api.Test;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import java.time.LocalDate;
import java.time.ZoneId;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
/**
* {@link RetentionAggregateServiceImpl} D1 次留口径单测 Mockito把守质量模型 §3.4 口径
* 首玩落点D1 窗口首玩日+1玩家×游戏×日去重匿名/登录稳定标识best-effort 零风险
*
* @author 绘境AI
*/
class RetentionAggregateServiceImplTest extends BaseMockitoUnitTest {
@InjectMocks
private RetentionAggregateServiceImpl service;
@Mock
private PlayerFirstPlayMapper playerFirstPlayMapper;
@Mock
private RetentionStatMapper retentionStatMapper;
private static final LocalDate D = LocalDate.of(2026, 7, 2);
/** 把 LocalDate 转成落在该自然日内的 ts平台时区 = systemDefault与实现同口径取正午防边界。 */
private static long tsForDate(LocalDate date) {
return date.atTime(12, 0).atZone(ZoneId.systemDefault()).toInstant().toEpochMilli();
}
private static EnvelopeReqVO playStart(String userId, String anonId, String gameId, LocalDate date) {
EnvelopeReqVO env = new EnvelopeReqVO();
env.setEventId("evt-" + userId + "-" + date);
env.setEvent("game_play_start");
env.setSchemaVersion("v1");
env.setTs(tsForDate(date));
env.setTraceId("t-" + date);
EnvelopeReqVO.UserVO user = new EnvelopeReqVO.UserVO();
user.setUserId(userId);
user.setAnonId(anonId);
env.setUser(user);
EnvelopeReqVO.ContextVO ctx = new EnvelopeReqVO.ContextVO();
ctx.setGameId(gameId);
env.setContext(ctx);
return env;
}
/** 首玩INSERT IGNORE 真插入(=1) → 累加 cohort 新玩家数、不置 D1。 */
@Test
void testFirstPlay_countsNewPlayer() {
when(playerFirstPlayMapper.insertIgnoreFirstPlay(eq(1024L), eq("u:10086"), eq(D))).thenReturn(1);
service.recordIfPlayStart(playStart("10086", null, "1024", D), 10086L);
verify(retentionStatMapper).insertOrAccumulateNewPlayer(eq(1024L), eq(D));
verify(retentionStatMapper, never()).insertOrAccumulateD1Retained(anyLong(), any());
verify(playerFirstPlayMapper, never()).markD1ReturnedIfUnset(anyLong());
}
/** D1 回访:首玩日+1 开局,未置位 → CAS 置位成功 → 累加分子cohort 日 = 首玩日 D。 */
@Test
void testD1Return_countsRetained() {
when(playerFirstPlayMapper.insertIgnoreFirstPlay(eq(1024L), eq("u:10086"), eq(D.plusDays(1)))).thenReturn(0);
PlayerFirstPlayDO row = new PlayerFirstPlayDO();
row.setId(7L);
row.setGameId(1024L);
row.setPlayerKey("u:10086");
row.setFirstPlayDate(D);
row.setD1Returned(0);
when(playerFirstPlayMapper.selectByGameAndPlayer(1024L, "u:10086")).thenReturn(row);
when(playerFirstPlayMapper.markD1ReturnedIfUnset(7L)).thenReturn(1);
service.recordIfPlayStart(playStart("10086", null, "1024", D.plusDays(1)), 10086L);
// 分子累加落在 cohort 首玩日 D非回访日 D+1
verify(retentionStatMapper).insertOrAccumulateD1Retained(eq(1024L), eq(D));
verify(retentionStatMapper, never()).insertOrAccumulateNewPlayer(anyLong(), any());
}
/** 同日重复开局:非 D1 窗口 → 不置位、不累加分子(去重)。 */
@Test
void testSameDayRepeat_noRetained() {
when(playerFirstPlayMapper.insertIgnoreFirstPlay(eq(1024L), eq("u:10086"), eq(D))).thenReturn(0);
PlayerFirstPlayDO row = new PlayerFirstPlayDO();
row.setId(7L);
row.setFirstPlayDate(D);
row.setD1Returned(0);
when(playerFirstPlayMapper.selectByGameAndPlayer(1024L, "u:10086")).thenReturn(row);
service.recordIfPlayStart(playStart("10086", null, "1024", D), 10086L);
verify(playerFirstPlayMapper, never()).markD1ReturnedIfUnset(anyLong());
verify(retentionStatMapper, never()).insertOrAccumulateD1Retained(anyLong(), any());
}
/** D+2 回访:超出 D1 窗口(首玩日+1→ 不计 D1。 */
@Test
void testD2Return_notCountedAsD1() {
when(playerFirstPlayMapper.insertIgnoreFirstPlay(eq(1024L), eq("u:10086"), eq(D.plusDays(2)))).thenReturn(0);
PlayerFirstPlayDO row = new PlayerFirstPlayDO();
row.setId(7L);
row.setFirstPlayDate(D);
row.setD1Returned(0);
when(playerFirstPlayMapper.selectByGameAndPlayer(1024L, "u:10086")).thenReturn(row);
service.recordIfPlayStart(playStart("10086", null, "1024", D.plusDays(2)), 10086L);
verify(playerFirstPlayMapper, never()).markD1ReturnedIfUnset(anyLong());
verify(retentionStatMapper, never()).insertOrAccumulateD1Retained(anyLong(), any());
}
/** 已置位d1_returned=1D1 窗口内重复回访不再累加CAS 去重的语义前置守卫)。 */
@Test
void testAlreadyReturned_noDoubleCount() {
when(playerFirstPlayMapper.insertIgnoreFirstPlay(eq(1024L), eq("u:10086"), eq(D.plusDays(1)))).thenReturn(0);
PlayerFirstPlayDO row = new PlayerFirstPlayDO();
row.setId(7L);
row.setFirstPlayDate(D);
row.setD1Returned(1); // 已置位
when(playerFirstPlayMapper.selectByGameAndPlayer(1024L, "u:10086")).thenReturn(row);
service.recordIfPlayStart(playStart("10086", null, "1024", D.plusDays(1)), 10086L);
verify(playerFirstPlayMapper, never()).markD1ReturnedIfUnset(anyLong());
verify(retentionStatMapper, never()).insertOrAccumulateD1Retained(anyLong(), any());
}
/** 未登录:无 userId、有 anonId → 玩家标识 a:{anonId}。 */
@Test
void testAnonymousPlayer_usesAnonIdKey() {
when(playerFirstPlayMapper.insertIgnoreFirstPlay(eq(1024L), eq("a:anon-xyz"), eq(D))).thenReturn(1);
service.recordIfPlayStart(playStart(null, "anon-xyz", "1024", D), null);
verify(playerFirstPlayMapper).insertIgnoreFirstPlay(eq(1024L), eq("a:anon-xyz"), eq(D));
verify(retentionStatMapper).insertOrAccumulateNewPlayer(eq(1024L), eq(D));
}
/** 非 game_play_start 事件 → no-op不碰任何 mapper。 */
@Test
void testNonPlayStart_noOp() {
EnvelopeReqVO like = playStart("10086", null, "1024", D);
like.setEvent("like");
service.recordIfPlayStart(like, 10086L);
verifyNoInteractions(playerFirstPlayMapper);
verifyNoInteractions(retentionStatMapper);
}
/** 无 gameId / 无玩家标识 → no-op无法归 cohort。 */
@Test
void testNoGameIdOrNoPlayer_noOp() {
service.recordIfPlayStart(playStart("10086", null, null, D), 10086L); // gameId
service.recordIfPlayStart(playStart(null, null, "1024", D), null); // userId 也无 anonId
verifyNoInteractions(playerFirstPlayMapper);
verifyNoInteractions(retentionStatMapper);
}
/** best-effortmapper 抛异常 → recordIfPlayStart 不外抛(保 quality_score 主链零风险)。 */
@Test
void testBestEffort_swallowsMapperException() {
when(playerFirstPlayMapper.insertIgnoreFirstPlay(any(), any(), any()))
.thenThrow(new RuntimeException("DB down"));
assertDoesNotThrow(() -> service.recordIfPlayStart(playStart("10086", null, "1024", D), 10086L));
}
}

View File

@ -73,4 +73,12 @@ public interface ErrorCodeConstants {
/** 订阅赋值幂等冲突(同一 bizNo 并发提交,唯一键 uk_biz_no 冲突且回查不到——非幂等的并发异常) */
ErrorCode TRADE_SUBSCRIPTION_BIZ_NO_CONFLICT = new ErrorCode(1_106_006_002, "订阅赋值提交冲突,请重试");
// ========== 积分充值 1-106-007-***P0 骨架mock-gated 收单真实支付收单 P1 受日历闸门==========
/** 充值金额/积分非法amount ≤ 0 或 points ≤ 0充值只增不减必须为正 */
ErrorCode TRADE_CHARGE_AMOUNT_INVALID = new ErrorCode(1_106_007_000, "充值金额与积分必须大于 0");
/** 充值幂等冲突(同一 bizNo 并发提交,唯一键 uk_biz_no 冲突且回查不到——非幂等的并发异常) */
ErrorCode TRADE_CHARGE_BIZ_NO_CONFLICT = new ErrorCode(1_106_007_001, "充值提交冲突,请重试");
/** 充值渠道不可用Nacos trade.charge-channel 配的渠道未注册实现,如真实 wxpay/alipay 未接入fail-fast 受日历闸门,绝不静默降级 mock 发假单) */
ErrorCode TRADE_CHARGE_CHANNEL_UNAVAILABLE = new ErrorCode(1_106_007_002, "充值渠道不可用,请稍后重试");
}

View File

@ -0,0 +1,58 @@
package com.wanxiang.huijing.game.module.trade.controller.app;
import com.wanxiang.huijing.game.module.trade.controller.app.vo.ChargeCreateReqVO;
import com.wanxiang.huijing.game.module.trade.controller.app.vo.ChargeOrderRespVO;
import com.wanxiang.huijing.game.module.trade.dal.dataobject.ChargeOrderDO;
import com.wanxiang.huijing.game.module.trade.service.charge.ChargeService;
import com.wanxiang.huijing.framework.common.pojo.CommonResult;
import com.wanxiang.huijing.framework.common.util.object.BeanUtils;
import com.wanxiang.huijing.framework.security.core.util.SecurityFrameworkUtils;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.annotation.Resource;
import jakarta.validation.Valid;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import static com.wanxiang.huijing.framework.common.pojo.CommonResult.success;
/**
* 产品端game-studio- 积分充值控制器P0 骨架
*
* <p>端前缀 /app-api 由框架按包名 controller.app.* 自动添加鉴权需用户 TokenuserId token 解析
* 用户只能给自己充值只见自己积分数据边界在 Service 强制真实支付收单 P1 受日历闸门
* MVP mock 渠道同步直通即入账Nacos trade.charge-channel=mock真实渠道未接入时 Service 侧工厂 fail-fast
*
* @author 绘境AI
*/
@Tag(name = "用户 App - 积分充值")
@RestController
@RequestMapping("/trade/charge")
@Validated
public class AppChargeController {
@Resource
private ChargeService chargeService;
@PostMapping("/create")
@Operation(summary = "创建积分充值单", description = "mock 渠道同步直通:创建即支付成功并积分入账;幂等以 bizNo 去重")
public CommonResult<ChargeOrderRespVO> createCharge(@Valid @RequestBody ChargeCreateReqVO reqVO) {
// userId token 解析用户只能给自己充值不接受前端传入 userId防越权
Long userId = SecurityFrameworkUtils.getLoginUserId();
ChargeOrderDO order = chargeService.createCharge(
userId, reqVO.getAmount(), reqVO.getPoints(), reqVO.getBizNo(), reqVO.getRemark());
return success(BeanUtils.toBean(order, ChargeOrderRespVO.class));
}
@GetMapping("/points-balance")
@Operation(summary = "查我的积分余额", description = "= Σ 本人已支付充值单的 points账本自包含消耗留后续")
public CommonResult<Long> getPointsBalance() {
Long userId = SecurityFrameworkUtils.getLoginUserId();
return success(chargeService.getPointsBalance(userId));
}
}

View File

@ -0,0 +1,35 @@
package com.wanxiang.huijing.game.module.trade.controller.app.vo;
import io.swagger.v3.oas.annotations.media.Schema;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Positive;
import lombok.Data;
/**
* 积分充值创建 Request VOP0 骨架
*
* @author 绘境AI
*/
@Schema(description = "产品端 - 积分充值创建 Request VO")
@Data
public class ChargeCreateReqVO {
@Schema(description = "充值金额(单位:分,>0", requiredMode = Schema.RequiredMode.REQUIRED, example = "1000")
@NotNull(message = "充值金额不能为空")
@Positive(message = "充值金额必须大于 0")
private Long amount;
@Schema(description = "充值获得积分(单位:点,>0", requiredMode = Schema.RequiredMode.REQUIRED, example = "100")
@NotNull(message = "充值积分不能为空")
@Positive(message = "充值积分必须大于 0")
private Long points;
@Schema(description = "充值业务单号 = 幂等键(客户端生成唯一)", requiredMode = Schema.RequiredMode.REQUIRED, example = "charge-20260702-abc")
@NotBlank(message = "充值业务单号不能为空")
private String bizNo;
@Schema(description = "备注(充值说明,可空)", example = "充值 100 积分")
private String remark;
}

View File

@ -0,0 +1,38 @@
package com.wanxiang.huijing.game.module.trade.controller.app.vo;
import io.swagger.v3.oas.annotations.media.Schema;
import lombok.Data;
import java.time.LocalDateTime;
/**
* 积分充值单 Response VOP0 骨架
*
* @author 绘境AI
*/
@Schema(description = "产品端 - 积分充值单 Response VO")
@Data
public class ChargeOrderRespVO {
@Schema(description = "充值单 ID", example = "1024")
private Long id;
@Schema(description = "充值金额(单位:分)", example = "1000")
private Long amount;
@Schema(description = "充值获得积分(单位:点)", example = "100")
private Long points;
@Schema(description = "充值渠道码mock/wxpay/alipay", example = "mock")
private String channel;
@Schema(description = "状态0待支付 10已支付 20已取消/失败", example = "10")
private Integer status;
@Schema(description = "渠道支付单号(供对账)", example = "mock-charge-1024-charge-20260702-abc")
private String payRef;
@Schema(description = "支付成功时间(已支付时有值)")
private LocalDateTime payTime;
}

View File

@ -0,0 +1,47 @@
package com.wanxiang.huijing.game.module.trade.dal.dataobject;
import com.wanxiang.huijing.framework.tenant.core.db.TenantBaseDO;
import com.baomidou.mybatisplus.annotation.KeySequence;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import lombok.EqualsAndHashCode;
import java.time.LocalDateTime;
/**
* 积分充值单 DO对应表 game_trade_charge_order积分充值 P0 骨架
*
* <p>继承 {@link TenantBaseDO} 自动携带 creator/create_time/updater/update_time/deleted + tenant_id
* 状态机0待支付 10已支付(积分入账) / 20已取消或失败金额/积分一律/BIGINT 禁浮点
* 幂等uk_biz_no用户积分余额 = Σ 本人 status=10 points账本自包含不污染 game_trade_account 不变式
*
* @author 绘境AI
*/
@TableName("game_trade_charge_order")
@KeySequence("game_trade_charge_order_seq")
@Data
@EqualsAndHashCode(callSuper = true)
public class ChargeOrderDO extends TenantBaseDO {
/** 充值单 ID */
private Long id;
/** 充值用户 IDDataPermission用户只见自己的充值单 */
private Long userId;
/** 充值金额(单位:分,>0 */
private Long amount;
/** 充值获得积分(单位:点,>0 */
private Long points;
/** 充值渠道码Nacos trade.charge-channelmock/wxpay/alipay */
private String channel;
/** 状态机0待支付 10已支付(积分入账) 20已取消/失败 */
private Integer status;
/** 充值业务单号 = 幂等键uk_biz_no */
private String bizNo;
/** 渠道支付单号mock 为本地 mock 单号;真实渠道为收单单号,供对账) */
private String payRef;
/** 支付成功时间status 置 10 时写入) */
private LocalDateTime payTime;
/** 备注(充值说明,可空) */
private String remark;
}

View File

@ -0,0 +1,60 @@
package com.wanxiang.huijing.game.module.trade.dal.mysql;
import com.wanxiang.huijing.game.module.trade.dal.dataobject.ChargeOrderDO;
import com.wanxiang.huijing.framework.mybatis.core.mapper.BaseMapperX;
import com.wanxiang.huijing.framework.mybatis.core.query.LambdaQueryWrapperX;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Select;
import org.apache.ibatis.annotations.Update;
import java.time.LocalDateTime;
/**
* 积分充值单 Mapper收单骨架
*
* <p>幂等靠 {@link #selectByBizNo}uk_biz_no 回查已支付置位靠 {@link #markPaidCas}status 010 CAS恰一次入账
* 用户积分余额靠 {@link #sumPaidPoints}Σ status=10 points账本自包含
*
* @author 绘境AI
*/
@Mapper
public interface ChargeOrderMapper extends BaseMapperX<ChargeOrderDO> {
/**
* 按业务单号取充值单幂等回查同一 bizNo 命中即返回已有单不重复收单
*
* @param bizNo 充值业务单号
* @return 充值单不存在返回 null
*/
default ChargeOrderDO selectByBizNo(String bizNo) {
return selectOne(new LambdaQueryWrapperX<ChargeOrderDO>()
.eq(ChargeOrderDO::getBizNo, bizNo));
}
/**
* CAS 置已支付status 010恰一次入账去重仅当当前为 0待支付 才置 10已支付 + pay_ref/pay_time
*
* <p>返回受影响行数保证并发/重放下同一单只入账一次1=本次置位成功调用方据此视为入账完成
* 0=已被置位或非待支付态幂等跳过
*
* @param id 充值单 ID
* @param payRef 渠道支付单号
* @param payTime 支付成功时间
* @return 受影响行数1=首次置位0=已置位/非待支付
*/
@Update("UPDATE game_trade_charge_order SET status = 10, pay_ref = #{payRef}, pay_time = #{payTime} "
+ "WHERE id = #{id} AND status = 0")
int markPaidCas(@Param("id") Long id, @Param("payRef") String payRef, @Param("payTime") LocalDateTime payTime);
/**
* 用户积分余额 = Σ 本人已支付(status=10)充值单的 points账本自包含不含消耗消耗留后续
*
* @param userId 用户 ID
* @return 积分余额无已支付单返回 0
*/
@Select("SELECT IFNULL(SUM(points), 0) FROM game_trade_charge_order "
+ "WHERE user_id = #{userId} AND status = 10 AND deleted = 0")
long sumPaidPoints(@Param("userId") Long userId);
}

View File

@ -0,0 +1,43 @@
package com.wanxiang.huijing.game.module.trade.framework.charge;
/**
* 充值收单渠道 SPIChargeClient统一形态屏蔽 mock/wxpay/alipay 差异业务码只认渠道码不认具体收单钱包
*
* <p>对应 payout {@code PayoutClient} 的镜像范式复用 huijing/yudao pay 模式最小改动
* ChargeService 只调本 SPI {@link #initiateCharge} 发起收单不在业务层散落微信/支付宝分支
* 渠道由 {@link ChargeClientFactory} Nacos {@code trade.charge-channel} 选择默认 mock切渠道=改配置业务码不变
*
* <p><b>mock-gated受日历闸门</b>MVP 不接真实支付收单trade.yaml支付收单=P1
* mock 渠道 = 同步直通发起即支付成功本进程直接入账真实 wxpay/alipay 渠道未接入时工厂 fail-fast
* {@code TRADE_CHARGE_CHANNEL_UNAVAILABLE}绝不静默降级 mock 收假单 payout假打款吞真钱红线同构收单侧=收假钱
*
* @author 绘境AI
*/
public interface ChargeClient {
/**
* 本实现对应的渠道码mock/wxpay/alipay Nacos trade.charge-channel 取值一致
*
* @return 渠道码
*/
String getChannel();
/**
* 是否同步直通渠道发起即支付成功本进程直接入账
*
* <p>mock=true同步直通无外部异步回调wxpay/alipay=false等收单 notify 回调驱动终态
* ChargeService 据此决定同步渠道发起后本进程直接置已支付 + 入账异步渠道仅发起等外部回调MVP 不涉及
*
* @return true=同步直通false=异步回调
*/
boolean isSynchronous();
/**
* 发起收单同步渠道发起即成功异步渠道仅发起终态等回调
*
* @param req 收单发起请求充值单 id / 金额 / 积分 / 用户 / 业务单号
* @return 发起受理结果含渠道支付单号 payRef发起失败抛 ServiceException调用方 fail-fast不静默降级
*/
ChargeResult initiateCharge(ChargeRequest req);
}

View File

@ -0,0 +1,74 @@
package com.wanxiang.huijing.game.module.trade.framework.charge;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import java.util.List;
import java.util.Map;
import java.util.function.Function;
import java.util.stream.Collectors;
import static com.wanxiang.huijing.game.module.trade.enums.ErrorCodeConstants.TRADE_CHARGE_CHANNEL_UNAVAILABLE;
import static com.wanxiang.huijing.framework.common.exception.util.ServiceExceptionUtil.exception;
/**
* ChargeClient 工厂 Nacos trade.charge-channel 路由到对应充值收单渠道 SPI 实现镜像 PayoutClientFactory
*
* <p>Spring 注入全部 {@link ChargeClient} 实现 getChannel() 建索引渠道选择走 Nacostrade.charge-channel默认 mock
* ChargeService 只问工厂要当前渠道客户端不认具体收单钱包切渠道=改配置业务码零分支
*
* <p><b>降级红线 payout 同构</b>配置渠道未命中实现时 <b>fail-fast 抛错挂起</b>
* {@code TRADE_CHARGE_CHANNEL_UNAVAILABLE}<b>严禁静默兜底到 mock</b>真实 wxpay/alipay 未接入时收假单
* = 收假钱/给假积分必须挂起受日历闸门支付资质下证前真实渠道不注册mock {@link MockChargeClient} 提供MVP 默认
*
* @author 绘境AI
*/
@Component
public class ChargeClientFactory {
private static final Logger log = LoggerFactory.getLogger(ChargeClientFactory.class);
/**
* 当前充值渠道走 Nacostrade.charge-channel默认 mockMVP
* @Value 占位符默认值仅为Nacos 未配置时的兜底线上以 Nacos 为准真实化时改为 wxpay/alipay 即切业务码零分支
*/
@Value("${trade.charge-channel:mock}")
private String chargeChannel;
/** 渠道码 → 实现,应用启动时一次性构建 */
private final Map<String, ChargeClient> clientMap;
public ChargeClientFactory(List<ChargeClient> clients) {
this.clientMap = clients.stream()
.collect(Collectors.toMap(ChargeClient::getChannel, Function.identity()));
}
/**
* 取当前 Nacos 配置渠道的充值收单客户端
*
* <p>红线配置渠道未注册实现 fail-fast {@code TRADE_CHARGE_CHANNEL_UNAVAILABLE}严禁静默降级 mock 收假单
*
* @return 当前渠道 ChargeClient保证非 null否则已抛错
*/
public ChargeClient current() {
ChargeClient client = clientMap.get(chargeChannel);
if (client == null) {
log.error("[ChargeClientFactory] 充值渠道不可用 channel={} 已注册渠道={}fail-fast不降级 mock",
chargeChannel, clientMap.keySet());
throw exception(TRADE_CHARGE_CHANNEL_UNAVAILABLE);
}
return client;
}
/**
* 当前配置的渠道码 charge_order.channel
*
* @return 渠道码来自 Nacos trade.charge-channel
*/
public String currentChannel() {
return chargeChannel;
}
}

View File

@ -0,0 +1,26 @@
package com.wanxiang.huijing.game.module.trade.framework.charge;
import lombok.AllArgsConstructor;
import lombok.Data;
/**
* 充值收单发起请求充值渠道 SPI 入参屏蔽 mock/wxpay/alipay 差异
*
* @author 绘境AI
*/
@Data
@AllArgsConstructor
public class ChargeRequest {
/** 充值单 ID已落库 game_trade_charge_order.id */
private Long chargeId;
/** 充值业务单号(幂等键,落 pay_ref 拼接可对账) */
private String bizNo;
/** 充值用户 ID */
private Long userId;
/** 充值金额(单位:分,>0 */
private Long amount;
/** 充值获得积分(单位:点,>0 */
private Long points;
}

View File

@ -0,0 +1,23 @@
package com.wanxiang.huijing.game.module.trade.framework.charge;
import lombok.Data;
/**
* 充值收单发起受理结果充值渠道 SPI 出参
*
* <p>发起受理结果含渠道支付单号 {@code payRef}mock 同步直通渠道发起即支付成功
* 真实异步渠道wxpay/alipay发起后终态由 notify 回调驱动MVP 未接入受日历闸门
*
* @author 绘境AI
*/
@Data
public class ChargeResult {
/** 渠道支付单号mock 为本地 mock 单号;真实渠道为收单单号,落 pay_ref 供对账) */
private final String payRef;
public ChargeResult(String payRef) {
this.payRef = payRef;
}
}

View File

@ -0,0 +1,45 @@
package com.wanxiang.huijing.game.module.trade.framework.charge;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
/**
* mock 充值收单渠道实现MVP 默认 trade.charge-channel=mock
*
* <p>同步直通发起收单即受理成功返回本地 mock 支付单号终态不在本类置 ChargeService 在发起后本进程
* 同步驱动 status 010已支付+ 积分入账保持与异步渠道一致的状态机路径不抄近路直接置 10
*
* <p>真实渠道wxpay/alipay走支付资质下证 + pay 进件后由独立 ChargeClient 实现注入本类不变工厂按 Nacos
* trade.charge-channel 路由mock 为默认兜底真实渠道缺失 fail-fast不静默降级 mock 收假单
*
* @author 绘境AI
*/
@Component
public class MockChargeClient implements ChargeClient {
private static final Logger log = LoggerFactory.getLogger(MockChargeClient.class);
/** 渠道码mock */
public static final String CHANNEL = "mock";
@Override
public String getChannel() {
return CHANNEL;
}
@Override
public boolean isSynchronous() {
return true; // mock 同步直通发起后本进程直接驱动已支付 + 入账
}
@Override
public ChargeResult initiateCharge(ChargeRequest req) {
// mock 支付单号以充值单 id + bizNo 拼接保证同一充值单稳定可对账非真实渠道单号
String mockRef = "mock-charge-" + req.getChargeId() + "-" + req.getBizNo();
log.info("[MockChargeClient] mock 发起收单同步直通chargeId={} bizNo={} amount={} points={} payRef={}",
req.getChargeId(), req.getBizNo(), req.getAmount(), req.getPoints(), mockRef);
return new ChargeResult(mockRef);
}
}

View File

@ -0,0 +1,36 @@
package com.wanxiang.huijing.game.module.trade.service.charge;
import com.wanxiang.huijing.game.module.trade.dal.dataobject.ChargeOrderDO;
/**
* 积分充值 ServiceP0 骨架mock-gated 收单
*
* <p>创建充值单 经充值渠道 SPI{@code ChargeClientFactory}Nacos trade.charge-channel默认 mock发起收单
* mock 同步直通即置已支付 + 积分入账真实 wxpay/alipay 渠道未接入时工厂 fail-fast受日历闸门不收假单
* 幂等 bizNo 命中 uk_biz_no 去重用户积分余额 = Σ 本人已支付充值单 points账本自包含
*
* @author 绘境AI
*/
public interface ChargeService {
/**
* 创建并支付充值单幂等mock 渠道同步直通创建即支付成功 + 积分入账
*
* @param userId 充值用户 ID
* @param amount 充值金额单位>0
* @param points 充值获得积分单位>0
* @param bizNo 充值业务单号 = 幂等键调用端生成唯一
* @param remark 备注可空
* @return 充值单mock 同步渠道下为已支付态幂等命中返回已有单
*/
ChargeOrderDO createCharge(Long userId, Long amount, Long points, String bizNo, String remark);
/**
* 用户积分余额 = Σ 本人已支付(status=10)充值单 points账本自包含不含消耗消耗留后续
*
* @param userId 用户 ID
* @return 积分余额无已支付单返回 0
*/
long getPointsBalance(Long userId);
}

View File

@ -0,0 +1,119 @@
package com.wanxiang.huijing.game.module.trade.service.charge;
import com.wanxiang.huijing.game.module.trade.dal.dataobject.ChargeOrderDO;
import com.wanxiang.huijing.game.module.trade.dal.mysql.ChargeOrderMapper;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeClient;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeClientFactory;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeRequest;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeResult;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.util.StringUtils;
import java.time.LocalDateTime;
import static com.wanxiang.huijing.game.module.trade.enums.ErrorCodeConstants.TRADE_CHARGE_AMOUNT_INVALID;
import static com.wanxiang.huijing.game.module.trade.enums.ErrorCodeConstants.TRADE_CHARGE_BIZ_NO_CONFLICT;
import static com.wanxiang.huijing.framework.common.exception.util.ServiceExceptionUtil.exception;
/**
* 积分充值 Service 实现P0 骨架mock-gated 收单
*
* <p>状态机常量0待支付 10已支付(入账) / 20已取消mock 渠道同步直通创建后本进程直接驱动 010CAS
* 资金安全与可对账同 payout 侧红线真实渠道未接入时工厂 fail-fast不静默降级 mock 收假单幂等以 bizNo 去重
*
* @author 绘境AI
*/
@Slf4j
@Service
public class ChargeServiceImpl implements ChargeService {
/** 状态:待支付(创建初值) */
private static final int STATUS_CREATED = 0;
/** 状态:已支付(积分入账) */
static final int STATUS_PAID = 10;
@Resource
private ChargeOrderMapper chargeOrderMapper;
@Resource
private ChargeClientFactory chargeClientFactory;
@Override
@Transactional(rollbackFor = Exception.class) // 创建+发起收单+同步入账在同一本地事务任一步失败整单回滚不留脏单
public ChargeOrderDO createCharge(Long userId, Long amount, Long points, String bizNo, String remark) {
// 参数校验充值只增不减金额与积分必须 > 0
if (amount == null || amount <= 0 || points == null || points <= 0) {
throw exception(TRADE_CHARGE_AMOUNT_INVALID);
}
// 幂等回查同一 bizNo 命中即返回已有单不重复收单不重复入账
ChargeOrderDO existing = chargeOrderMapper.selectByBizNo(bizNo);
if (existing != null) {
log.info("[charge] 幂等命中,返回已有充值单 bizNo={} chargeId={} status={}", bizNo, existing.getId(), existing.getStatus());
return existing;
}
// 取当前渠道客户端Nacos trade.charge-channel默认 mock真实渠道未注册 fail-fast绝不静默降级 mock 收假单
ChargeClient client = chargeClientFactory.current();
String channel = chargeClientFactory.currentChannel();
// 创建充值单status=待支付uk_biz_no 兜底并发冲突则回查返回已有单保证幂等
ChargeOrderDO order = new ChargeOrderDO();
order.setUserId(userId);
order.setAmount(amount);
order.setPoints(points);
order.setChannel(channel);
order.setStatus(STATUS_CREATED);
order.setBizNo(bizNo);
order.setPayRef("");
order.setRemark(StringUtils.hasText(remark) ? remark : "");
try {
chargeOrderMapper.insert(order);
} catch (DuplicateKeyException dup) {
// 并发同 bizNo另一线程已插入 回查返回已有单幂等回查不到则为异常并发抛冲突
ChargeOrderDO concurrent = chargeOrderMapper.selectByBizNo(bizNo);
if (concurrent != null) {
log.info("[charge] 并发同 bizNo返回已插入充值单 bizNo={} chargeId={}", bizNo, concurrent.getId());
return concurrent;
}
log.error("[charge] 充值单 uk_biz_no 冲突但回查不到 bizNo={}", bizNo, dup);
throw exception(TRADE_CHARGE_BIZ_NO_CONFLICT);
}
// 发起收单外部交互留痕mock 同步直通返回 mock 支付单号
ChargeResult result = client.initiateCharge(new ChargeRequest(order.getId(), bizNo, userId, amount, points));
// 同步渠道mock发起即支付成功 CAS 置已支付 + 积分入账账本自包含已支付单即入账凭据
if (client.isSynchronous()) {
LocalDateTime payTime = LocalDateTime.now();
int paid = chargeOrderMapper.markPaidCas(order.getId(), result.getPayRef(), payTime);
if (paid > 0) {
order.setStatus(STATUS_PAID);
order.setPayRef(result.getPayRef());
order.setPayTime(payTime);
log.info("[charge] mock 同步充值成功入账 chargeId={} userId={} amount={} points={} payRef={}",
order.getId(), userId, amount, points, result.getPayRef());
} else {
// 极端并发CAS 未命中已被置位回读当前态返回不重复入账
ChargeOrderDO refreshed = chargeOrderMapper.selectById(order.getId());
if (refreshed != null) {
return refreshed;
}
}
} else {
// 异步渠道wxpay/alipayMVP 未接入仅发起 pay_ref终态等收单 notify 回调驱动本骨架不涉及
order.setPayRef(result.getPayRef());
log.info("[charge] 异步渠道发起收单,待 notify 驱动终态 chargeId={} channel={} payRef={}",
order.getId(), channel, result.getPayRef());
}
return order;
}
@Override
public long getPointsBalance(Long userId) {
return chargeOrderMapper.sumPaidPoints(userId);
}
}

View File

@ -0,0 +1,52 @@
package com.wanxiang.huijing.game.module.trade.framework.charge;
import com.wanxiang.huijing.framework.common.exception.ServiceException;
import org.junit.jupiter.api.Test;
import org.springframework.test.util.ReflectionTestUtils;
import java.util.List;
import static com.wanxiang.huijing.game.module.trade.enums.ErrorCodeConstants.TRADE_CHARGE_CHANNEL_UNAVAILABLE;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* 充值渠道 SPI 单元测试MockChargeClient 契约 + ChargeClientFactory 路由与 fail-fast 红线镜像 payout
*
* @author 绘境AI
*/
class ChargeChannelTest {
/** MockChargeClient 契约channel=mock、同步直通、payRef 稳定可对账(含 chargeId+bizNo。 */
@Test
void testMockChargeClient_contract() {
MockChargeClient client = new MockChargeClient();
assertEquals("mock", client.getChannel());
assertTrue(client.isSynchronous(), "mock 应同步直通");
ChargeResult r = client.initiateCharge(new ChargeRequest(1024L, "biz-1", 100L, 1000L, 100L));
assertEquals("mock-charge-1024-biz-1", r.getPayRef(), "mock 支付单号应含 chargeId+bizNo 稳定可对账");
}
/** 工厂路由Nacos trade.charge-channel=mock → 命中 MockChargeClient。 */
@Test
void testFactory_current_routesToMock() {
MockChargeClient mock = new MockChargeClient();
ChargeClientFactory factory = new ChargeClientFactory(List.of(mock));
ReflectionTestUtils.setField(factory, "chargeChannel", "mock");
assertSame(mock, factory.current());
assertEquals("mock", factory.currentChannel());
}
/** fail-fast 红线:配置真实渠道(wxpay)但未注册实现 → 抛 TRADE_CHARGE_CHANNEL_UNAVAILABLE绝不静默降级 mock 收假单。 */
@Test
void testFactory_unregisteredChannel_failFast() {
ChargeClientFactory factory = new ChargeClientFactory(List.of(new MockChargeClient()));
ReflectionTestUtils.setField(factory, "chargeChannel", "wxpay"); // 真实渠道未接入受日历闸门
ServiceException ex = assertThrows(ServiceException.class, factory::current);
assertEquals(TRADE_CHARGE_CHANNEL_UNAVAILABLE.getCode(), ex.getCode(), "真实渠道缺失应 fail-fast 不降级 mock");
}
}

View File

@ -0,0 +1,150 @@
package com.wanxiang.huijing.game.module.trade.service.charge;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.SQLIntegrityConstraintViolationException;
import java.sql.Statement;
import java.sql.Timestamp;
import java.time.LocalDateTime;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* R2charge mock e2e集成证据 mini-infra MySQL{@code 100.64.0.8:3306}·Tailscale 直连·绕系统 SOCKS 代理
* throwaway 库里跑<b> V28 迁移文件的 DDL</b> + 积分充值 mock 收单流程 SQL验入账/幂等/积分余额正确
*
* <p><b>覆盖边界</b> MySQL V28 DDL 可跑双验 Flyway DDL 合法性 charge 骨架 SQLinsert 待支付单 /
* markPaidCas 010 入账 / sumPaidPoints 汇总积分余额 / uk_biz_no 幂等去重行为正确ChargeService 分支逻辑
* 校验/幂等/mock 同步入账 {@link ChargeServiceImplTest} 单测坐实 IT 用真 SQL 复演其收单序列
* 全程 throwaway {@code trade_r2_charge_e2e} /隔离不碰真库真表不与 Flyway 冲突
*
* <p><b>默认跳过</b>{@code @EnabledIfSystemProperty(trade.charge.e2e=1)}常规 {@code mvn test} 不连 MySQL
* 出证据时 {@code mvn test -Dtrade.charge.e2e=1 -Dtest=ChargeOrderRealDbIT}
*
* @author 绘境AI
*/
@EnabledIfSystemProperty(named = "trade.charge.e2e", matches = "1")
class ChargeOrderRealDbIT {
private static final String HOST = System.getProperty("trade.charge.dbhost", "100.64.0.8");
private static final String USER = System.getProperty("trade.charge.dbuser", "root");
// 内网 MVP 凭据铁律授权入仓 docs/内网凭据与端点.md可经 -D 覆盖
private static final String PASS = System.getProperty("trade.charge.dbpass", "ZRH3jwYLOrntBcTAw29MW9BP");
private static final String DB = "trade_r2_charge_e2e";
private static final long UID = 999_999_002L;
private String baseUrl() {
return "jdbc:mysql://" + HOST + ":3306/?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=UTC";
}
@Test
void v28Ddl_and_chargeMockFlow_onRealMysql() throws Exception {
// 内网直连绕系统 SOCKS 代理macOS SOCKSEnable127.0.0.1:7897MySQL Connector/J 默认走 JVM socks须显式禁
System.setProperty("socksProxyHost", "");
System.setProperty("java.net.useSystemProxies", "false");
try (Connection conn = DriverManager.getConnection(baseUrl(), USER, PASS)) {
try (Statement st = conn.createStatement()) {
st.execute("DROP DATABASE IF EXISTS " + DB);
st.execute("CREATE DATABASE " + DB + " DEFAULT CHARACTER SET utf8mb4");
st.execute("USE " + DB);
}
runV28Ddl(conn);
// 充值单 1创建待支付 mock 同步置已支付入账points=100
long id1 = insertCreated(conn, "biz-e2e-1", 1000L, 100L);
int paid1 = markPaidCas(conn, id1, "mock-charge-" + id1 + "-biz-e2e-1");
assertEquals(1, paid1, "首次 CAS 0→10 应成功");
// 重复入账幂等 markPaidCas 同单 0 已非待支付
assertEquals(0, markPaidCas(conn, id1, "dup-ref"), "已支付单重复 CAS 应 0 行(不重复入账)");
// 充值单 2另一 bizNo入账 points=50
long id2 = insertCreated(conn, "biz-e2e-2", 500L, 50L);
assertEquals(1, markPaidCas(conn, id2, "mock-charge-" + id2 + "-biz-e2e-2"));
// 积分余额 = Σ 已支付单 points = 100 + 50 = 150
assertEquals(150L, sumPaidPoints(conn), "积分余额应为已支付单 points 之和");
// bizNo 幂等重复插入同 biz-e2e-1 uk_biz_no 唯一键冲突收单不重复
assertThrows(SQLIntegrityConstraintViolationException.class,
() -> insertCreated(conn, "biz-e2e-1", 1000L, 100L), "同 bizNo 应被 uk_biz_no 拒绝");
} finally {
try (Connection conn = DriverManager.getConnection(baseUrl(), USER, PASS);
Statement st = conn.createStatement()) {
st.execute("DROP DATABASE IF EXISTS " + DB);
}
}
}
private void runV28Ddl(Connection conn) throws Exception {
Path migration = Paths.get(System.getProperty("user.dir"), "..", "..",
"huijing-server", "src", "main", "resources", "db", "migration",
"V28.0.0__create_game_trade_charge_order.sql");
assertTrue(Files.exists(migration), "应找到 V28 迁移文件:" + migration.toAbsolutePath().normalize());
String sql = Files.readString(migration);
StringBuilder cleaned = new StringBuilder();
for (String line : sql.split("\n")) {
if (line.trim().startsWith("--")) {
continue;
}
cleaned.append(line).append('\n');
}
try (Statement st = conn.createStatement()) {
for (String stmt : cleaned.toString().split(";")) {
if (!stmt.trim().isEmpty()) {
st.execute(stmt);
}
}
}
}
/** 复演 ChargeService插入待支付单status=0。返回自增 id。uk_biz_no 冲突抛 SQLIntegrityConstraintViolationException。 */
private long insertCreated(Connection conn, String bizNo, long amount, long points) throws SQLException {
try (PreparedStatement ps = conn.prepareStatement(
"INSERT INTO game_trade_charge_order (user_id, amount, points, channel, status, biz_no, pay_ref) "
+ "VALUES (?,?,?, 'mock', 0, ?, '')", Statement.RETURN_GENERATED_KEYS)) {
ps.setLong(1, UID);
ps.setLong(2, amount);
ps.setLong(3, points);
ps.setString(4, bizNo);
ps.executeUpdate();
try (ResultSet rs = ps.getGeneratedKeys()) {
rs.next();
return rs.getLong(1);
}
}
}
/** 复演 ChargeOrderMapper.markPaidCasstatus 0→10 CAS + 落 pay_ref/pay_time。返回受影响行数。 */
private int markPaidCas(Connection conn, long id, String payRef) throws SQLException {
try (PreparedStatement ps = conn.prepareStatement(
"UPDATE game_trade_charge_order SET status = 10, pay_ref = ?, pay_time = ? WHERE id = ? AND status = 0")) {
ps.setString(1, payRef);
ps.setTimestamp(2, Timestamp.valueOf(LocalDateTime.now()));
ps.setLong(3, id);
return ps.executeUpdate();
}
}
/** 复演 ChargeOrderMapper.sumPaidPointsΣ status=10 的 points。 */
private long sumPaidPoints(Connection conn) throws SQLException {
try (PreparedStatement ps = conn.prepareStatement(
"SELECT IFNULL(SUM(points), 0) FROM game_trade_charge_order WHERE user_id = ? AND status = 10 AND deleted = 0")) {
ps.setLong(1, UID);
try (ResultSet rs = ps.executeQuery()) {
rs.next();
return rs.getLong(1);
}
}
}
}

View File

@ -0,0 +1,106 @@
package com.wanxiang.huijing.game.module.trade.service.charge;
import com.wanxiang.huijing.game.module.trade.dal.dataobject.ChargeOrderDO;
import com.wanxiang.huijing.game.module.trade.dal.mysql.ChargeOrderMapper;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeClient;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeClientFactory;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeRequest;
import com.wanxiang.huijing.game.module.trade.framework.charge.ChargeResult;
import com.wanxiang.huijing.framework.common.exception.ServiceException;
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
import org.junit.jupiter.api.Test;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import java.time.LocalDateTime;
import static com.wanxiang.huijing.game.module.trade.enums.ErrorCodeConstants.TRADE_CHARGE_AMOUNT_INVALID;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
/**
* {@link ChargeServiceImpl} 单元测试 Mockito把守积分充值 mock-gated 收单流程幂等参数校验
*
* @author 绘境AI
*/
class ChargeServiceImplTest extends BaseMockitoUnitTest {
@InjectMocks
private ChargeServiceImpl chargeService;
@Mock
private ChargeOrderMapper chargeOrderMapper;
@Mock
private ChargeClientFactory chargeClientFactory;
@Mock
private ChargeClient mockChargeClient;
/** mock 同步充值创建→发起收单→CAS 置已支付+积分入账单据落已支付态、payRef 落库。 */
@Test
void testCreateCharge_mockSynchronous_paidAndCredited() {
when(chargeOrderMapper.selectByBizNo("biz-1")).thenReturn(null); // 非幂等命中
when(chargeClientFactory.current()).thenReturn(mockChargeClient);
when(chargeClientFactory.currentChannel()).thenReturn("mock");
when(mockChargeClient.isSynchronous()).thenReturn(true);
when(mockChargeClient.initiateCharge(any(ChargeRequest.class))).thenReturn(new ChargeResult("mock-ref-1"));
when(chargeOrderMapper.markPaidCas(any(), eq("mock-ref-1"), any(LocalDateTime.class))).thenReturn(1);
ChargeOrderDO order = chargeService.createCharge(100L, 1000L, 100L, "biz-1", "充值 100 积分");
// 落库为待支付初值 CAS 置已支付
verify(chargeOrderMapper).insert(any(ChargeOrderDO.class));
verify(mockChargeClient).initiateCharge(any(ChargeRequest.class));
verify(chargeOrderMapper).markPaidCas(any(), eq("mock-ref-1"), any(LocalDateTime.class));
assertEquals(ChargeServiceImpl.STATUS_PAID, order.getStatus(), "mock 同步应置已支付");
assertEquals("mock-ref-1", order.getPayRef());
assertEquals("mock", order.getChannel());
assertEquals(100L, order.getPoints());
}
/** 幂等:同 bizNo 命中已有单 → 返回已有单,不重复插入/不重复发起收单/不重复入账。 */
@Test
void testCreateCharge_idempotentOnBizNo() {
ChargeOrderDO existing = new ChargeOrderDO();
existing.setId(7L);
existing.setBizNo("biz-1");
existing.setStatus(ChargeServiceImpl.STATUS_PAID);
when(chargeOrderMapper.selectByBizNo("biz-1")).thenReturn(existing);
ChargeOrderDO result = chargeService.createCharge(100L, 1000L, 100L, "biz-1", null);
assertEquals(7L, result.getId());
verify(chargeOrderMapper, never()).insert(any(ChargeOrderDO.class));
verify(chargeClientFactory, never()).current();
}
/** 参数校验:金额 ≤ 0 → 抛 TRADE_CHARGE_AMOUNT_INVALID不落库、不发起收单。 */
@Test
void testCreateCharge_invalidAmount_throws() {
ServiceException ex = assertThrows(ServiceException.class,
() -> chargeService.createCharge(100L, 0L, 100L, "biz-1", null));
assertEquals(TRADE_CHARGE_AMOUNT_INVALID.getCode(), ex.getCode());
verify(chargeOrderMapper, never()).insert(any(ChargeOrderDO.class));
}
/** 参数校验:积分 ≤ 0 → 抛 TRADE_CHARGE_AMOUNT_INVALID。 */
@Test
void testCreateCharge_invalidPoints_throws() {
ServiceException ex = assertThrows(ServiceException.class,
() -> chargeService.createCharge(100L, 1000L, 0L, "biz-1", null));
assertEquals(TRADE_CHARGE_AMOUNT_INVALID.getCode(), ex.getCode());
}
/** 积分余额 = Σ 本人已支付单 points委托 mapper.sumPaidPoints。 */
@Test
void testGetPointsBalance_sumsPaidOrders() {
when(chargeOrderMapper.sumPaidPoints(100L)).thenReturn(250L);
assertEquals(250L, chargeService.getPointsBalance(100L));
verify(chargeOrderMapper).sumPaidPoints(anyLong());
}
}

View File

@ -190,6 +190,13 @@ rocketmq:
producer:
group: ${spring.application.name}_PRODUCER # 生产者分组
# telemetry 遥测异步(切片一 阶段〇 R1批量事件入口改 MQ 异步消费的灰度开关
telemetry:
mq:
# true=上报入口快速 ACK + 投递 telemetry-event 主题、消费侧异步聚合(幂等);
# false默认=生产者/消费者 Bean 缺席、回落原请求线程内同步写(现行行为逐字不变)。经 .env 注入按环境开启。
enabled: ${TELEMETRY_MQ_ENABLED:false}
spring:
# Kafka 配置项,对应 KafkaProperties 配置类
kafka:

View File

@ -0,0 +1,59 @@
-- =============================================================================
-- 契约 #2 DB 迁移 | 模块telemetryD1 次留旁路聚合)| ownerWS5
-- 文件V27.0.0__create_game_telemetry_retention.sqlFlyway只新增已合入禁止修改回滚写新补偿迁移
-- 内容D1 次留旁路两表 —— 玩家首玩追踪 game_telemetry_player_first_play + 游戏维度次留聚合 game_telemetry_retention_stat
-- 背景:质量模型设计 §3.4 L3 留存结构立到「可建字段」。既有 game_telemetry_game_stat 结构与 quality_score 链路
-- 零改动、零风险——次留是跨日玩家级去重 + cohort 追踪,不是 game_stat 那种当日行级累加,故独立旁路表承载。
-- D1 口径§3.4):同一 game_id、同一稳定玩家标识登录取 userId、未登录取匿名设备标识 anonId对齐 #5 契约 anonId 口径)
-- 在首玩自然日(平台时区,与 game_telemetry_game_stat.stat_date 同日界)之后的第 1 个自然日内 ≥1 次开局事件game_play_start
-- 按「玩家 × 游戏 × 日」去重。D1 留存率 = d1_retained_count / new_player_count读时计算、不落列避免除法陈旧
-- 增量落法:消费侧对 game_play_start 事件旁路累加best-effort失败不连累 quality_score 主链),见 RetentionAggregateService。
-- 约定InnoDB + utf8mb4显式列含 Yudao 审计列 + 租户列。
-- 错误码段telemetry = 1-104-***-***
-- =============================================================================
-- -----------------------------------------------------------------------------
-- 表game_telemetry_player_first_play —— 玩家首玩追踪(游戏 × 玩家 粒度)
-- 每玩家对每游戏一行:记首玩自然日 + 是否在首玩日+1 自然日回访开局d1_returned CAS 置位一次去重)。
-- 玩家标识命名空间隔离:登录=u:{userId}、未登录=a:{anonId},防 userId 与 anonId 字面碰撞。
-- 已知 MVP 边界:匿名→登录跨身份(先 anon 玩、次日登录玩当前记为两个玩家user_login 事件桥接留后续精化。
-- -----------------------------------------------------------------------------
CREATE TABLE `game_telemetry_player_first_play` (
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '记录 ID',
`game_id` BIGINT NOT NULL COMMENT '游戏 ID= game_project.id',
`player_key` VARCHAR(80) NOT NULL COMMENT '稳定玩家标识:登录=u:{userId},未登录=a:{anonId}(对齐 #5 契约 anonId 口径;两类命名空间隔离防碰撞)',
`first_play_date` DATE NOT NULL COMMENT '该玩家对该游戏的首玩自然日(平台时区,与 game_telemetry_game_stat.stat_date 同日界)',
`d1_returned` TINYINT NOT NULL DEFAULT 0 COMMENT '次留标记0未回 1已在首玩日+1 自然日回访开局CAS 置位一次、去重)',
`creator` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '创建者Yudao 审计列;上报通道为 system/上报方)',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updater` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '更新者',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
`deleted` BIT(1) NOT NULL DEFAULT b'0' COMMENT '逻辑删除0未删 1已删',
`tenant_id` BIGINT NOT NULL DEFAULT 0 COMMENT '租户 IDYudao 多租户兼容MVP 单租户=0',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_game_player` (`game_id`, `player_key`) COMMENT '每玩家每游戏一行(首玩落点 + D1 去重的幂等键)',
KEY `idx_game_firstdate` (`game_id`, `first_play_date`) COMMENT '按游戏+首玩日取 cohort对账/回算)'
) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4 COMMENT = '遥测玩家首玩追踪表D1 次留旁路,游戏×玩家粒度)';
-- -----------------------------------------------------------------------------
-- 表game_telemetry_retention_stat —— 游戏维度 D1 次留聚合(游戏 × cohort日 粒度)
-- 由消费侧对 game_play_start 增量累加产出;唯一约束 (game_id, stat_date) 保证聚合幂等(行级原子累加)。
-- new_player_countcohort 规模当日首玩新玩家数d1_retained_count该 cohort 次日回访开局玩家数。
-- 独立于 game_telemetry_game_stat不动既有聚合表结构、quality_score 链路零风险;本表暂不回灌 feed观测/回流用)。
-- -----------------------------------------------------------------------------
CREATE TABLE `game_telemetry_retention_stat` (
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '聚合记录 ID',
`game_id` BIGINT NOT NULL COMMENT '游戏 ID= game_project.id',
`stat_date` DATE NOT NULL COMMENT 'cohort 日 = 首玩自然日(平台时区)',
`new_player_count` BIGINT NOT NULL DEFAULT 0 COMMENT '当日首玩新玩家数cohort 规模,玩家×游戏×日去重)',
`d1_retained_count` BIGINT NOT NULL DEFAULT 0 COMMENT '当日 cohort 中在次日回访开局的玩家数D1 留存分子D1 率 = 本列/new_player_count 读时算)',
`creator` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '创建者Yudao 审计列;聚合任务为 system',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updater` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '更新者',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
`deleted` BIT(1) NOT NULL DEFAULT b'0' COMMENT '逻辑删除0未删 1已删',
`tenant_id` BIGINT NOT NULL DEFAULT 0 COMMENT '租户 IDYudao 多租户兼容MVP 单租户=0',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_game_date` (`game_id`, `stat_date`) COMMENT '聚合幂等游戏×cohort日一行消费侧行级原子累加',
KEY `idx_date` (`stat_date`) COMMENT '按日取留存(看板/数据飞轮回流)'
) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4 COMMENT = '游戏维度 D1 次留聚合表(旁路,不动既有 game_telemetry_game_stat';

View File

@ -0,0 +1,38 @@
-- =============================================================================
-- 契约 #2 DB 迁移 | 模块trade积分充值 P0 骨架)| ownerWS5
-- 文件V28.0.0__create_game_trade_charge_order.sqlFlyway只新增已合入禁止修改回滚写新补偿迁移
-- 内容:积分充值单表 game_trade_charge_order —— 用户充值获得积分的收单骨架
-- 背景:补 55 P0 最后一块结构缺口。MVP 不接真实支付收单trade.yaml会员订阅/支付收单=P1、归 pay 日历闸门)——
-- 充值渠道复用 trade 既有 payout 那套「mock-gated 渠道 SPI」范式Nacos trade.charge-channel默认 mock路由
-- mock 同步直通即支付成功入账、真实 wxpay/alipay 渠道未接入时工厂 fail-fast受日历闸门绝不静默发假单
-- 账本自包含:用户积分余额 = Σ 本人已支付(status=10)充值单的 points不污染创作者收益账户(game_trade_account)不变式。
-- 约定InnoDB + utf8mb4金额/积分一律「分/点」BIGINT 禁浮点;含 Yudao 审计列 + 租户列。
-- 错误码段trade = 1-106-***-***charge 用 007 子段)
-- =============================================================================
-- -----------------------------------------------------------------------------
-- 表game_trade_charge_order —— 积分充值单(收单骨架)
-- 状态机0待支付 → 10已支付(积分入账) / 20已取消或失败。mock 渠道同步直通:创建后本进程直接驱动 0→10。
-- 幂等uk_biz_no —— 同一 biz_no 重复提交返回已有单(不重复收单、不重复入账)。
-- -----------------------------------------------------------------------------
CREATE TABLE `game_trade_charge_order` (
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '充值单 ID',
`user_id` BIGINT NOT NULL COMMENT '充值用户 IDDataPermission用户只见自己的充值单',
`amount` BIGINT NOT NULL COMMENT '充值金额单位BIGINT 禁浮点,>0',
`points` BIGINT NOT NULL COMMENT '充值获得积分(单位:点,>0积分余额 = Σ 本人已支付单 points',
`channel` VARCHAR(16) NOT NULL DEFAULT 'mock' COMMENT '充值渠道码Nacos trade.charge-channelmock/wxpay/alipay',
`status` TINYINT NOT NULL DEFAULT 0 COMMENT '状态机0待支付 10已支付(积分入账) 20已取消/失败',
`biz_no` VARCHAR(64) NOT NULL COMMENT '充值业务单号 = 幂等键(调用端生成唯一)',
`pay_ref` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '渠道支付单号mock 为本地 mock 单号;真实渠道为收单单号,供对账)',
`pay_time` DATETIME NULL COMMENT '支付成功时间status 置 10 时写入)',
`remark` VARCHAR(255) NOT NULL DEFAULT '' COMMENT '备注(充值说明,可空)',
`creator` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '创建者Yudao 审计列)',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updater` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '更新者',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
`deleted` BIT(1) NOT NULL DEFAULT b'0' COMMENT '逻辑删除0未删 1已删',
`tenant_id` BIGINT NOT NULL DEFAULT 0 COMMENT '租户 IDYudao 多租户兼容MVP 单租户=0',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_biz_no` (`biz_no`) COMMENT '充值幂等键:同一 biz_no 只收单一次',
KEY `idx_user_status` (`user_id`, `status`) COMMENT '按用户取充值单 / 汇总积分余额(Σ status=10 的 points'
) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4 COMMENT = '积分充值单收单骨架mock-gated真实支付收单 P1 受日历闸门)';