fix(community): U3 评审整改——P0 ZSET 移出事务(afterCommit) + P1 弹幕长度门/弹幕广播出事务/举报接缝去重 + P2 注释与公开读
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
3f35acb4d3
commit
edc3436f7e
@ -14,6 +14,7 @@ import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import jakarta.annotation.security.PermitAll;
|
||||
import jakarta.validation.Valid;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
@ -55,6 +56,7 @@ public class AppCommunityCommentController {
|
||||
}
|
||||
|
||||
@GetMapping("/comments")
|
||||
@PermitAll // 公开读:评论列表为公开内容,不取 userId(按 target 维度,匿名可浏览),放行匿名对齐文档
|
||||
@Operation(summary = "评论分页列表", description = "公开内容,按 target_type+target_id 维度,id 倒序")
|
||||
public CommonResult<PageResult<CommentRespVO>> getCommentPage(@Valid CommentPageReqVO pageReqVO) {
|
||||
PageResult<CommentDO> pageResult = commentService.getCommentPage(pageReqVO);
|
||||
@ -62,6 +64,7 @@ public class AppCommunityCommentController {
|
||||
}
|
||||
|
||||
@GetMapping("/comments/count")
|
||||
@PermitAll // 公开读:评论计数为公开内容,不取 userId(匿名可读),放行匿名对齐文档
|
||||
@Operation(summary = "评论数", description = "某目标的评论总数(公开计数)")
|
||||
@Parameter(name = "targetType", description = "目标类型:1游戏", example = "1")
|
||||
@Parameter(name = "targetId", description = "目标 ID(gameId)", required = true, example = "1024")
|
||||
|
||||
@ -9,6 +9,7 @@ import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import jakarta.annotation.security.PermitAll;
|
||||
import jakarta.validation.Valid;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
@ -44,6 +45,7 @@ public class AppCommunityDanmakuController {
|
||||
}
|
||||
|
||||
@GetMapping("/danmaku/recent")
|
||||
@PermitAll // 公开读:弹幕回放为公开内容,不取 userId(进房补帧,匿名可读),放行匿名对齐文档
|
||||
@Operation(summary = "最近弹幕回放", description = "公开读;某游戏最近 N 条弹幕(进房补帧),按时间正序")
|
||||
@Parameter(name = "gameId", description = "游戏编号(房间)", required = true, example = "1024")
|
||||
@Parameter(name = "limit", description = "取近 N 条(默认 50)", example = "50")
|
||||
|
||||
@ -9,6 +9,7 @@ import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import jakarta.annotation.security.PermitAll;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
|
||||
@ -49,6 +50,7 @@ public class AppCommunityFollowController {
|
||||
}
|
||||
|
||||
@GetMapping("/follow/stat/{userId}")
|
||||
@PermitAll // 公开读:粉丝/关注数为公开内容;getLoginUserId() 已 null 安全(未登录则 followedByMe 恒 false),放行匿名对齐文档
|
||||
@Operation(summary = "关注统计", description = "某用户粉丝数/关注数 + 当前登录者是否已关注 ta(未登录则 followedByMe=false)")
|
||||
@Parameter(name = "userId", description = "查询的用户编号", required = true, example = "88")
|
||||
public CommonResult<FollowStatRespVO> getFollowStat(@PathVariable("userId") Long userId) {
|
||||
|
||||
@ -7,6 +7,7 @@ import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import jakarta.annotation.security.PermitAll;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
@ -35,6 +36,7 @@ public class AppCommunityRankController {
|
||||
private RankService rankService;
|
||||
|
||||
@GetMapping("/rank")
|
||||
@PermitAll // 公开读:榜单为公开内容,无 userId 依赖(匿名玩家可浏览),放行匿名以对齐文档「无需登录」
|
||||
@Operation(summary = "排行榜单 TopN", description = "公开读;维度须为已定义维度(game_hot/creator_hot);ZSET 缺失自动 DB 回填")
|
||||
@Parameter(name = "dimension", description = "排行维度:game_hot 游戏热度 / creator_hot 创作者热度", required = true, example = "game_hot")
|
||||
@Parameter(name = "topN", description = "取前 N 名(默认 10)", example = "10")
|
||||
|
||||
@ -110,10 +110,11 @@ public class CommentServiceImpl implements CommentService {
|
||||
update.setId(id);
|
||||
update.setReportStatus(1); // 1=被举报
|
||||
commentMapper.updateById(update);
|
||||
// 转 compliance 受理接缝须与状态机一致:仅在真正 0→1 跃迁时注册,避免对同一评论重复举报多次触发受理(幂等)。
|
||||
// 事务提交后转 compliance 受理(事务内禁止远程调用 §7)
|
||||
registerComplianceIntakeAfterCommit(id, comment.getTargetId(), userId, reason);
|
||||
}
|
||||
log.info("[reportComment] 评论举报已落本端举报态 commentId={}, reporterUserId={}, reason={}", id, userId, reason);
|
||||
// 事务提交后转 compliance 受理(事务内禁止远程调用 §7)
|
||||
registerComplianceIntakeAfterCommit(id, comment.getTargetId(), userId, reason);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@ -61,7 +61,12 @@ public class DanmakuRoomRegistry {
|
||||
}
|
||||
|
||||
/**
|
||||
* 会话从所有房间移除(连接断开时调用,防泄漏残留会话)
|
||||
* 会话从所有房间移除(防泄漏残留会话)。
|
||||
*
|
||||
* 现状:本波【未】接框架断连回调(afterConnectionClosed),即不在断连瞬间调用本方法;
|
||||
* 安静房间内的失效会话靠广播时的懒清理(见 DanmakuServiceImpl#broadcastToRoom,doSend 前剔除已关会话)兜底。
|
||||
* 本方法仍在用:退房 leave 与单测调用;保留以备断连即时清理接通。
|
||||
* OPEN ITEM (多实例/断连即时清理波): wire WS afterConnectionClosed → removeSession
|
||||
*
|
||||
* @param sessionId WebSocket 会话 ID
|
||||
*/
|
||||
|
||||
@ -10,7 +10,6 @@ import com.wanxiang.huijing.framework.websocket.core.session.WebSocketSessionMan
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
import org.springframework.web.socket.WebSocketSession;
|
||||
|
||||
import java.util.ArrayList;
|
||||
@ -19,6 +18,7 @@ import java.util.List;
|
||||
import java.util.Set;
|
||||
|
||||
import static com.wanxiang.huijing.game.module.community.enums.ErrorCodeConstants.COMMUNITY_DANMAKU_CONTENT_BLANK;
|
||||
import static com.wanxiang.huijing.game.module.community.enums.ErrorCodeConstants.COMMUNITY_DANMAKU_CONTENT_TOO_LONG;
|
||||
import static com.wanxiang.huijing.framework.common.exception.util.ServiceExceptionUtil.exception;
|
||||
|
||||
/**
|
||||
@ -41,6 +41,9 @@ public class DanmakuServiceImpl implements DanmakuService {
|
||||
/** WebSocket 弹幕消息类型(前后端约定,前端按此 type 渲染弹幕轨道) */
|
||||
public static final String WS_TYPE_DANMAKU = "danmaku";
|
||||
|
||||
/** 弹幕正文最大长度(与 DB 列 content VARCHAR(255) 对齐,服务端可信边界拦截,覆盖 HTTP/WS 两路) */
|
||||
private static final int CONTENT_MAX_LEN = 255;
|
||||
|
||||
@Resource
|
||||
private DanmakuMapper danmakuMapper;
|
||||
|
||||
@ -58,15 +61,27 @@ public class DanmakuServiceImpl implements DanmakuService {
|
||||
@Resource
|
||||
private AbstractWebSocketMessageSender webSocketMessageSender;
|
||||
|
||||
/**
|
||||
* 发送弹幕 = 落库(可回放)+ 房间广播。
|
||||
*
|
||||
* 不加 @Transactional:本方法实为单条 insert(已自动提交),无需声明式事务;且广播是网络 IO,
|
||||
* 红线 §7/§4.1 禁在 DB 开事务内做网络 IO。去掉事务后,insert 先持久化、广播自然在持久化之后发出,
|
||||
* 不会出现「广播了却随后回滚」的幻象弹幕(客户端先收到、DB 却没有的不一致)。
|
||||
* 内容长度在此服务端可信边界拦截(覆盖 HTTP 与 WS 两路):WS 入站经 DanmakuWebSocketMessageListener
|
||||
* 程序化构造 ReqVO,绕过 Controller 的 @Size(255),仅靠 Controller 校验会让超长内容直插 VARCHAR(255) 致截断/插入失败。
|
||||
*/
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public DanmakuRespVO sendDanmaku(DanmakuSendReqVO reqVO, Long userId) {
|
||||
// 内容非空校验(去首尾空白)
|
||||
String content = reqVO.getContent() == null ? "" : reqVO.getContent().trim();
|
||||
if (content.isEmpty()) {
|
||||
throw exception(COMMUNITY_DANMAKU_CONTENT_BLANK);
|
||||
}
|
||||
// 落库(可回放)
|
||||
// 内容超长拦截(服务端可信边界,覆盖 HTTP/WS;DB 列为 VARCHAR(255),超长会截断/插入失败)
|
||||
if (content.length() > CONTENT_MAX_LEN) {
|
||||
throw exception(COMMUNITY_DANMAKU_CONTENT_TOO_LONG);
|
||||
}
|
||||
// 落库(可回放)——单条 insert 自动提交,持久化后再广播
|
||||
DanmakuDO danmaku = new DanmakuDO();
|
||||
danmaku.setGameId(reqVO.getGameId());
|
||||
danmaku.setUserId(userId); // 发送者=当前登录用户(token 解析,非前端入参)
|
||||
|
||||
@ -11,6 +11,8 @@ import org.springframework.data.redis.core.RedisTemplate;
|
||||
import org.springframework.data.redis.core.ZSetOperations;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
import org.springframework.transaction.support.TransactionSynchronization;
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
@ -24,10 +26,11 @@ import static com.wanxiang.huijing.framework.common.exception.util.ServiceExcept
|
||||
*
|
||||
* ============================ ZSET↔DB 一致性设计(核心) ============================
|
||||
* 角色:DB game_community_rank = 真相(source of truth);Redis ZSET community:rank:{dimension} = 缓存/索引。
|
||||
* 1) 写穿透(addScore):分值变更在同一 Service 流程内【双写】——
|
||||
* 先 DB 原子累加 score(@Transactional 保证 DB 侧;首次 insert 兜底,并发双插由 uk 裁决),
|
||||
* 再 ZSET incrementScore 同步索引。DB 提交后 ZSET 已写,二者最终一致。
|
||||
* 2) 一致性窗口(最终一致):DB 与 ZSET 在「DB 提交 → ZSET 写入」之间存在极短窗口不一致;
|
||||
* 1) 写穿透(addScore):分值变更【双写】——
|
||||
* DB 原子累加 score 在 @Transactional 内(首次 insert 兜底,并发双插由 uk 裁决),保证 DB 侧原子;
|
||||
* ZSET incrementScore 经 afterCommit 接缝在【DB 提交后】才执行(红线 §7/§4.1:Redis 非事务参与者,
|
||||
* 禁在 DB 开事务内做 Redis 写)。提交成功才写 ZSET,二者最终一致;DB 回滚则 ZSET 不写,无脏索引。
|
||||
* 2) 一致性窗口(最终一致):DB 与 ZSET 在「DB 提交 → afterCommit 写 ZSET」之间存在极短窗口不一致;
|
||||
* 若进程在两步之间崩溃,ZSET 会落后于 DB —— 由 (3) 的回填对账修复(DB 权威,ZSET 向 DB 收敛)。
|
||||
* 3) 回填重建(rebuildFromDb):ZSET 缺失(缓存丢失/未预热)/启动/周期对账时,以 DB TopN 重写 ZSET,
|
||||
* 消除漂移。getTopN 读 ZSET 为空时自动触发回填兜底,保证读路径不空窗。
|
||||
@ -57,15 +60,40 @@ public class RankServiceImpl implements RankService {
|
||||
public void addScore(String dimension, Long targetId, long delta) {
|
||||
validateDimension(dimension);
|
||||
// —— 第一写(权威):DB 原子累加;记录不存在则 insert 初始化兜底(并发双插由 uk 裁决再累加)——
|
||||
// 仅 DB 写在本事务内保持原子;Redis 写绝不入事务(见下)。
|
||||
int affected = rankMapper.incrScore(dimension, targetId, delta);
|
||||
if (affected == 0) {
|
||||
ensureRankRow(dimension, targetId, delta);
|
||||
}
|
||||
// —— 第二写(索引):ZSET 同步累加,与 DB 同向(最终一致)——
|
||||
// 说明:Redis 非事务参与者;此处在 DB 写之后同流程执行,DB 提交后 ZSET 即生效。
|
||||
// 若两步间崩溃致 ZSET 落后,由 rebuildFromDb 对账向 DB 收敛(DB 权威)。
|
||||
redisTemplate.opsForZSet().incrementScore(buildKey(dimension), String.valueOf(targetId), (double) delta);
|
||||
log.info("[addScore] 排行双写完成(DB 权威 + ZSET 索引)dimension={}, targetId={}, delta={}", dimension, targetId, delta);
|
||||
// —— 第二写(索引):ZSET 同步累加,与 DB 同向(最终一致),经 afterCommit 在 DB 提交后才执行 ——
|
||||
// 红线 §7/§4.1:Redis 非事务参与者,禁在 DB 开事务内做 Redis 写。改注册事务同步钩子,
|
||||
// DB 提交成功后才 incrementScore;DB 回滚则不写 ZSET(无脏索引)。
|
||||
// 若提交后崩溃致 ZSET 落后,由 rebuildFromDb 对账向 DB 收敛(DB 权威)。
|
||||
registerZsetIncrAfterCommit(dimension, targetId, delta);
|
||||
log.info("[addScore] 排行 DB 写入完成、ZSET 索引提交后写(DB 权威 + ZSET 索引)dimension={}, targetId={}, delta={}", dimension, targetId, delta);
|
||||
}
|
||||
|
||||
/**
|
||||
* 在事务提交后向 ZSET 累加分值(红线 §7/§4.1:Redis 写绝不入 DB 事务)。
|
||||
*
|
||||
* 用 afterCommit 而非事务内直写:① Redis 非事务参与者,事务内写无法随 DB 回滚而撤销,会留脏索引;
|
||||
* ② DB 提交才代表分值真正生效,此时再写 ZSET 才与 DB 权威同向。
|
||||
* 无活动事务(理论不至,@Transactional 已开)时直接写,避免静默丢索引更新(由 rebuildFromDb 兜底对账)。
|
||||
*/
|
||||
private void registerZsetIncrAfterCommit(String dimension, Long targetId, long delta) {
|
||||
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
|
||||
// 无活动事务:直接写 ZSET(兜底,避免丢索引);正常路径不会走到这里
|
||||
redisTemplate.opsForZSet().incrementScore(buildKey(dimension), String.valueOf(targetId), (double) delta);
|
||||
return;
|
||||
}
|
||||
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
|
||||
@Override
|
||||
public void afterCommit() {
|
||||
// DB 已提交,分值生效 → 此刻写 ZSET 索引,与 DB 同向(最终一致)
|
||||
redisTemplate.opsForZSet().incrementScore(buildKey(dimension), String.valueOf(targetId), (double) delta);
|
||||
log.info("[addScore][afterCommit] ZSET 索引提交后累加完成 dimension={}, targetId={}, delta={}", dimension, targetId, delta);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@ -5,10 +5,12 @@ import com.wanxiang.huijing.game.module.community.dal.dataobject.comment.Comment
|
||||
import com.wanxiang.huijing.game.module.community.dal.mysql.comment.CommentMapper;
|
||||
import com.wanxiang.huijing.framework.common.exception.ServiceException;
|
||||
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.ArgumentCaptor;
|
||||
import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
|
||||
import static com.wanxiang.huijing.game.module.community.enums.ErrorCodeConstants.*;
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@ -20,6 +22,8 @@ import static org.mockito.Mockito.*;
|
||||
*
|
||||
* 覆盖评论核心规则(U3 R-SOC):发表内容非空、删除归属校验(防越权删他人评论 — PII/数据边界红线)、
|
||||
* 举报理由必填 + 落举报态幂等 + 转 compliance 接缝。归属隔离 enforce 须有测试兜底。
|
||||
* compliance 接缝须与 0→1 跃迁一致:仅真正首次举报才注册 afterCommit 同步,重复举报不重复注册(幂等,
|
||||
* 用 initSynchronization() + 断言 getSynchronizations().size() 验证)。
|
||||
*
|
||||
* @author 造梦AI
|
||||
*/
|
||||
@ -120,27 +124,43 @@ class CommentServiceImplTest extends BaseMockitoUnitTest {
|
||||
void testReportComment_marksReportStatus() {
|
||||
CommentDO comment = ownedComment(1L, 11L, 0); // 正常态被举报
|
||||
when(commentMapper.selectByIdSafe(1L)).thenReturn(comment);
|
||||
// 模拟活动事务,观察 compliance 受理接缝是否在 0→1 跃迁时注册
|
||||
TransactionSynchronizationManager.initSynchronization();
|
||||
|
||||
commentService.reportComment(1L, "含违规", 99L);
|
||||
|
||||
ArgumentCaptor<CommentDO> captor = ArgumentCaptor.forClass(CommentDO.class);
|
||||
verify(commentMapper).updateById(captor.capture());
|
||||
assertEquals(1, captor.getValue().getReportStatus()); // 落举报态 1
|
||||
// 真正 0→1 跃迁:恰注册一个 afterCommit 接缝(转 compliance 受理)
|
||||
assertEquals(1, TransactionSynchronizationManager.getSynchronizations().size());
|
||||
}
|
||||
|
||||
@Test
|
||||
void testReportComment_alreadyReportedIdempotent() {
|
||||
CommentDO comment = ownedComment(1L, 11L, 1); // 已被举报
|
||||
when(commentMapper.selectByIdSafe(1L)).thenReturn(comment);
|
||||
// 模拟活动事务,验证重复举报不重复注册 compliance 接缝(幂等)
|
||||
TransactionSynchronizationManager.initSynchronization();
|
||||
|
||||
commentService.reportComment(1L, "含违规", 99L);
|
||||
|
||||
// 幂等:已举报态不重复写
|
||||
// 幂等:已举报态不重复写 DB
|
||||
verify(commentMapper, never()).updateById(any(CommentDO.class));
|
||||
// 幂等:非 0→1 跃迁不再注册 compliance 受理接缝(避免再次举报时 N 次触发受理)
|
||||
assertEquals(0, TransactionSynchronizationManager.getSynchronizations().size());
|
||||
}
|
||||
|
||||
// ============================== 测试夹具 ==============================
|
||||
|
||||
/** 每个用例后清理可能残留的事务同步上下文(仅 initSynchronization 后才需清,避免污染后续用例) */
|
||||
@AfterEach
|
||||
void tearDownSynchronization() {
|
||||
if (TransactionSynchronizationManager.isSynchronizationActive()) {
|
||||
TransactionSynchronizationManager.clearSynchronization();
|
||||
}
|
||||
}
|
||||
|
||||
/** 构造归属指定作者、指定举报态的评论 */
|
||||
private static CommentDO ownedComment(Long id, Long userId, Integer reportStatus) {
|
||||
CommentDO comment = new CommentDO();
|
||||
|
||||
@ -17,6 +17,7 @@ import java.util.List;
|
||||
import java.util.Set;
|
||||
|
||||
import static com.wanxiang.huijing.game.module.community.enums.ErrorCodeConstants.COMMUNITY_DANMAKU_CONTENT_BLANK;
|
||||
import static com.wanxiang.huijing.game.module.community.enums.ErrorCodeConstants.COMMUNITY_DANMAKU_CONTENT_TOO_LONG;
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyInt;
|
||||
@ -27,7 +28,8 @@ import static org.mockito.Mockito.*;
|
||||
/**
|
||||
* {@link DanmakuServiceImpl} 单元测试(纯 Mockito)
|
||||
*
|
||||
* 覆盖弹幕核心规则(U3 R-SOC):内容非空、落库 + 房间广播、空房间仅落库不广播、回放时间正序。
|
||||
* 覆盖弹幕核心规则(U3 R-SOC):内容非空、内容超长拦截(服务端可信边界,覆盖 WS 绕过 @Size 的路径)、
|
||||
* 落库 + 房间广播(落库后广播,无事务)、空房间仅落库不广播、回放时间正序。
|
||||
* 广播复用框架 doSend,验证从房间注册表取会话后下发。
|
||||
*
|
||||
* @author 造梦AI
|
||||
@ -60,6 +62,21 @@ class DanmakuServiceImplTest extends BaseMockitoUnitTest {
|
||||
verify(danmakuMapper, never()).insert(any(DanmakuDO.class));
|
||||
}
|
||||
|
||||
// ============================== 内容超长(服务端可信边界,覆盖 WS 绕过 @Size 路径) ==============================
|
||||
|
||||
@Test
|
||||
void testSendDanmaku_tooLongContentRejected() {
|
||||
// 256 字符(> VARCHAR(255)),模拟 WS 程序化构造绕过 Controller @Size(255) 的入参
|
||||
DanmakuSendReqVO reqVO = new DanmakuSendReqVO();
|
||||
reqVO.setGameId(1024L);
|
||||
reqVO.setContent("x".repeat(256));
|
||||
|
||||
ServiceException ex = assertThrows(ServiceException.class,
|
||||
() -> danmakuService.sendDanmaku(reqVO, 99L));
|
||||
assertEquals(COMMUNITY_DANMAKU_CONTENT_TOO_LONG.getCode(), ex.getCode());
|
||||
verify(danmakuMapper, never()).insert(any(DanmakuDO.class)); // 超长不落库
|
||||
}
|
||||
|
||||
// ============================== 落库 + 房间广播 ==============================
|
||||
|
||||
@Test
|
||||
|
||||
@ -5,6 +5,7 @@ import com.wanxiang.huijing.game.module.community.dal.dataobject.rank.RankDO;
|
||||
import com.wanxiang.huijing.game.module.community.dal.mysql.rank.RankMapper;
|
||||
import com.wanxiang.huijing.framework.common.exception.ServiceException;
|
||||
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
@ -12,6 +13,8 @@ import org.mockito.Mockito;
|
||||
import org.springframework.dao.DuplicateKeyException;
|
||||
import org.springframework.data.redis.core.RedisTemplate;
|
||||
import org.springframework.data.redis.core.ZSetOperations;
|
||||
import org.springframework.transaction.support.TransactionSynchronization;
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
@ -20,8 +23,10 @@ import java.util.Set;
|
||||
import static com.wanxiang.huijing.game.module.community.enums.ErrorCodeConstants.COMMUNITY_RANK_DIMENSION_INVALID;
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyDouble;
|
||||
import static org.mockito.ArgumentMatchers.anyInt;
|
||||
import static org.mockito.ArgumentMatchers.anyLong;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.eq;
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
@ -30,12 +35,14 @@ import static org.mockito.Mockito.*;
|
||||
*
|
||||
* 覆盖 ZSET↔DB 一致性核心逻辑(U3 R-SOC):
|
||||
* ① 维度合法性校验(脏维度拒);
|
||||
* ② 写穿透双写(DB 原子累加 + ZSET incrementScore 同流程都被调用);
|
||||
* ② 写穿透双写(DB 原子累加在事务内;ZSET incrementScore 延迟到 afterCommit 才执行——红线 §7/§4.1);
|
||||
* ③ 首次累加 0 行 → insert 初始化兜底(含并发 DuplicateKey CAS);
|
||||
* ④ getTopN 命中 ZSET 快路径;
|
||||
* ⑤ getTopN ZSET 缺失 → 触发 DB 回填重建后再读(读路径不空窗);
|
||||
* ⑥ rebuildFromDb 先删 key 再按 DB 权威回填。
|
||||
* ZSET 用 Mockito 桩,不依赖 jedismock ZSET 实现,专测一致性编排逻辑。
|
||||
* 事务接缝测法:addScore 前 initSynchronization() 模拟活动事务 → 调用后断言 ZSET 尚未写 →
|
||||
* 手动触发已注册同步的 afterCommit() → 再断言 ZSET 已写;teardown 调 clearSynchronization()。
|
||||
*
|
||||
* @author 造梦AI
|
||||
*/
|
||||
@ -64,46 +71,71 @@ class RankServiceImplTest extends BaseMockitoUnitTest {
|
||||
verify(rankMapper, never()).incrScore(any(), any(), anyLong());
|
||||
}
|
||||
|
||||
// ============================== 写穿透双写 ==============================
|
||||
// ============================== 写穿透双写(ZSET 延迟到 afterCommit) ==============================
|
||||
|
||||
@Test
|
||||
void testAddScore_writeThroughBothDbAndZset() {
|
||||
void testAddScore_zsetDeferredUntilAfterCommit() {
|
||||
// DB 原子累加命中已有行(affected=1,不走 insert)
|
||||
when(rankMapper.incrScore(DIM, 1024L, 5L)).thenReturn(1);
|
||||
when(redisTemplate.opsForZSet()).thenReturn(zSetOps);
|
||||
// 模拟存在活动事务(@Transactional 已开),以走 afterCommit 接缝
|
||||
TransactionSynchronizationManager.initSynchronization();
|
||||
|
||||
rankService.addScore(DIM, 1024L, 5L);
|
||||
|
||||
// DB 写在事务内已发生
|
||||
verify(rankMapper).incrScore(DIM, 1024L, 5L);
|
||||
verify(rankMapper, never()).insert(any(RankDO.class)); // 命中已有行不初始化
|
||||
// 关键:ZSET 写尚未发生(已注册同步,待提交后才执行)——红线 §7/§4.1
|
||||
verify(zSetOps, never()).incrementScore(anyString(), any(), anyDouble());
|
||||
assertEquals(1, TransactionSynchronizationManager.getSynchronizations().size());
|
||||
|
||||
// 触发提交后回调 → ZSET 此刻才同向累加
|
||||
fireAfterCommit();
|
||||
verify(zSetOps).incrementScore(ZKEY, "1024", 5.0d);
|
||||
}
|
||||
|
||||
@Test
|
||||
void testAddScore_firstTimeInsertFallback() {
|
||||
// 原子累加命中 0 行(记录不存在)→ insert 初始化兜底 → ZSET 提交后写
|
||||
when(rankMapper.incrScore(DIM, 1024L, 5L)).thenReturn(0);
|
||||
when(redisTemplate.opsForZSet()).thenReturn(zSetOps);
|
||||
TransactionSynchronizationManager.initSynchronization();
|
||||
|
||||
rankService.addScore(DIM, 1024L, 5L);
|
||||
|
||||
verify(rankMapper).insert(any(RankDO.class)); // 初始化兜底
|
||||
verify(zSetOps, never()).incrementScore(anyString(), any(), anyDouble()); // 提交前 ZSET 未写
|
||||
fireAfterCommit();
|
||||
verify(zSetOps).incrementScore(ZKEY, "1024", 5.0d); // 提交后 ZSET 同步
|
||||
}
|
||||
|
||||
@Test
|
||||
void testAddScore_concurrentInsertCasRetriesIncr() {
|
||||
// 累加 0 行 → insert 抛唯一键冲突(并发已建行)→ 补一次原子累加;ZSET 仍延迟到提交后
|
||||
when(rankMapper.incrScore(DIM, 1024L, 5L)).thenReturn(0).thenReturn(1);
|
||||
doThrow(new DuplicateKeyException("uk_dimension_target")).when(rankMapper).insert(any(RankDO.class));
|
||||
when(redisTemplate.opsForZSet()).thenReturn(zSetOps);
|
||||
TransactionSynchronizationManager.initSynchronization();
|
||||
|
||||
rankService.addScore(DIM, 1024L, 5L);
|
||||
|
||||
// 第一次 incr(0行) + catch 后补 incr:共 2 次(均在事务内)
|
||||
verify(rankMapper, times(2)).incrScore(DIM, 1024L, 5L);
|
||||
verify(zSetOps, never()).incrementScore(anyString(), any(), anyDouble()); // 提交前 ZSET 未写
|
||||
fireAfterCommit();
|
||||
verify(zSetOps).incrementScore(ZKEY, "1024", 5.0d);
|
||||
}
|
||||
|
||||
@Test
|
||||
void testAddScore_noActiveTransactionWritesZsetDirectly() {
|
||||
// 兜底分支:无活动事务(未 initSynchronization)→ 直接写 ZSET,避免丢索引
|
||||
when(rankMapper.incrScore(DIM, 1024L, 5L)).thenReturn(1);
|
||||
when(redisTemplate.opsForZSet()).thenReturn(zSetOps);
|
||||
|
||||
rankService.addScore(DIM, 1024L, 5L);
|
||||
|
||||
// 双写:DB 累加 + ZSET 同向累加
|
||||
verify(rankMapper).incrScore(DIM, 1024L, 5L);
|
||||
verify(zSetOps).incrementScore(ZKEY, "1024", 5.0d);
|
||||
verify(rankMapper, never()).insert(any(RankDO.class)); // 命中已有行不初始化
|
||||
}
|
||||
|
||||
@Test
|
||||
void testAddScore_firstTimeInsertFallback() {
|
||||
// 原子累加命中 0 行(记录不存在)→ insert 初始化兜底 → 仍写 ZSET
|
||||
when(rankMapper.incrScore(DIM, 1024L, 5L)).thenReturn(0);
|
||||
when(redisTemplate.opsForZSet()).thenReturn(zSetOps);
|
||||
|
||||
rankService.addScore(DIM, 1024L, 5L);
|
||||
|
||||
verify(rankMapper).insert(any(RankDO.class)); // 初始化兜底
|
||||
verify(zSetOps).incrementScore(ZKEY, "1024", 5.0d); // ZSET 仍同步
|
||||
}
|
||||
|
||||
@Test
|
||||
void testAddScore_concurrentInsertCasRetriesIncr() {
|
||||
// 累加 0 行 → insert 抛唯一键冲突(并发已建行)→ 补一次原子累加
|
||||
when(rankMapper.incrScore(DIM, 1024L, 5L)).thenReturn(0).thenReturn(1);
|
||||
doThrow(new DuplicateKeyException("uk_dimension_target")).when(rankMapper).insert(any(RankDO.class));
|
||||
when(redisTemplate.opsForZSet()).thenReturn(zSetOps);
|
||||
|
||||
rankService.addScore(DIM, 1024L, 5L);
|
||||
|
||||
// 第一次 incr(0行) + catch 后补 incr:共 2 次
|
||||
verify(rankMapper, times(2)).incrScore(DIM, 1024L, 5L);
|
||||
verify(zSetOps).incrementScore(ZKEY, "1024", 5.0d);
|
||||
verify(zSetOps).incrementScore(ZKEY, "1024", 5.0d); // 无事务时同流程直写
|
||||
}
|
||||
|
||||
// ============================== getTopN ==============================
|
||||
@ -173,6 +205,21 @@ class RankServiceImplTest extends BaseMockitoUnitTest {
|
||||
|
||||
// ============================== 测试夹具 ==============================
|
||||
|
||||
/** 每个用例后清理可能残留的事务同步上下文(仅在 initSynchronization 后才需清,避免污染后续用例) */
|
||||
@AfterEach
|
||||
void tearDownSynchronization() {
|
||||
if (TransactionSynchronizationManager.isSynchronizationActive()) {
|
||||
TransactionSynchronizationManager.clearSynchronization();
|
||||
}
|
||||
}
|
||||
|
||||
/** 手动触发当前已注册事务同步的 afterCommit()(模拟 DB 提交成功后框架回调) */
|
||||
private static void fireAfterCommit() {
|
||||
for (TransactionSynchronization sync : TransactionSynchronizationManager.getSynchronizations()) {
|
||||
sync.afterCommit();
|
||||
}
|
||||
}
|
||||
|
||||
/** 构造一个 ZSET 降序元组集合(1024→99 在前,2048→80 在后) */
|
||||
private static Set<ZSetOperations.TypedTuple<Object>> tuples() {
|
||||
Set<ZSetOperations.TypedTuple<Object>> set = new LinkedHashSet<>();
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user