package com.tailbet.job; import com.tailbet.model.entity.GameRound; import com.tailbet.model.entity.RedPacket; import com.tailbet.model.entity.User; import com.tailbet.mapper.GameRoundMapper; import com.tailbet.mapper.RedPacketMapper; import com.tailbet.mapper.RedPacketRecvMapper; import com.tailbet.mapper.UserMapper; import com.tailbet.mapper.DigitWhiteMapper; import com.tailbet.mapper.DrawRecordMapper; import com.tailbet.mapper.SysConfigMapper; import com.tailbet.model.entity.DigitWhite; import com.tailbet.model.entity.DrawRecord; import com.tailbet.model.entity.RedPacketRecv; import com.tailbet.model.entity.SysConfig; import com.tailbet.service.IGameService; import com.tailbet.service.IRedPacketService; import com.tailbet.service.IFakeUserService; import com.tailbet.openim.OpenIMClient; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.math.BigDecimal; import java.time.LocalDateTime; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; /** * 游戏定时任务 */ @Slf4j @Component @RequiredArgsConstructor public class GameJob { private final GameRoundMapper gameRoundMapper; private final RedPacketMapper redPacketMapper; private final RedPacketRecvMapper redPacketRecvMapper; private final DrawRecordMapper drawRecordMapper; private final SysConfigMapper sysConfigMapper; private final UserMapper userMapper; private final DigitWhiteMapper digitWhiteMapper; private final IGameService gameService; private final IRedPacketService redPacketService; private final IFakeUserService fakeUserService; private final OpenIMClient openIMClient; // 局状态监控(使用内存缓存,实际生产应使用Redis) private final Map roundTasks = new ConcurrentHashMap<>(); // 配置缓存 private final Map configCache = new ConcurrentHashMap<>(); // 延迟任务队列 private final Map pendingTasks = new ConcurrentHashMap<>(); // 任务ID生成器 private java.util.concurrent.atomic.AtomicLong taskIdGenerator = new java.util.concurrent.atomic.AtomicLong(1); /** * 启动时加载配置缓存 */ @jakarta.annotation.PostConstruct public void loadConfig() { try { List configs = sysConfigMapper.selectList(null); for (SysConfig config : configs) { configCache.put(config.getCfgKey(), config.getCfgValue()); } log.info("配置缓存加载完成,共 {} 条配置", configs.size()); } catch (Exception e) { log.error("配置缓存加载失败: {}", e.getMessage()); } } /** * 获取配置值 */ private String getConfig(String key) { return configCache.get(key); } /** * 获取整型配置值 */ public int getConfigInt(String key, int defaultValue) { String value = configCache.get(key); if (value == null) return defaultValue; try { return Integer.parseInt(value); } catch (NumberFormatException e) { return defaultValue; } } /** * 检查红包是否需要托领取(每3秒执行) */ @Scheduled(fixedDelay = 3000) public void checkRedPacketForFakeReceive() { // 查找进行中的游戏红包 List ongoingRps = redPacketMapper.selectList(new LambdaQueryWrapper() .eq(RedPacket::getStatus, 0) .eq(RedPacket::getType, IRedPacketService.TYPE_GAME) .isNotNull(RedPacket::getRoundId)); for (RedPacket rp : ongoingRps) { LocalDateTime createTime = rp.getCreateTime(); // 发出超过配置时间未领完,让托领取 int fakeRecvTimeout = getConfigInt("fake_recv_seconds", 5); if (createTime.plusSeconds(fakeRecvTimeout).isBefore(LocalDateTime.now())) { // 查找未领取的用户,随机选一个托领取 try { fakeUserReceiveRemaining(rp); } catch (Exception e) { log.error("托领取红包失败: rpId={}, error={}", rp.getId(), e.getMessage()); } } } } /** * 假人托领取剩余份额(循环领取直到红包领完) */ private void fakeUserReceiveRemaining(RedPacket rp) { List fakes = fakeUserService.getFakeUsers(); if (fakes.isEmpty()) { return; } // 循环领取直到红包领完 while (true) { RedPacket rpCheck = redPacketMapper.selectById(rp.getId()); if (rpCheck == null || rpCheck.getStatus() != 0) { // 红包已领完或不存在 break; } // 查找还未领取的假人 User fakeToUse = null; for (User fake : fakes) { // 检查该假人是否已领取过本红包 var received = redPacketRecvMapper.selectList( new LambdaQueryWrapper() .eq(RedPacketRecv::getRpId, rp.getId()) .eq(RedPacketRecv::getUserId, fake.getId())); if (received.isEmpty()) { fakeToUse = fake; break; } } if (fakeToUse == null) { // 所有假人都已领取过,无法继续 log.warn("所有假人都已领取过本红包,无法继续: rpId={}", rp.getId()); break; } try { redPacketService.receiveRp(rp.getId(), fakeToUse.getId()); log.info("假人{}领取红包{}剩余份额", fakeToUse.getId(), rp.getId()); } catch (Exception e) { log.error("假人领取红包失败: rpId={}, fakeId={}, error={}", rp.getId(), fakeToUse.getId(), e.getMessage()); // 可能是红包已领完,退出循环 break; } } } /** * 检查局超时(每10秒执行) */ @Scheduled(fixedDelay = 10000) public void checkRoundTimeout() { List ongoingRounds = gameRoundMapper.selectList(new LambdaQueryWrapper() .eq(GameRound::getStatus, 0)); for (GameRound round : ongoingRounds) { // 检查是否超过接注时间 RoundTask task = roundTasks.get(round.getId()); if (task != null && task.endBetTime.isBefore(LocalDateTime.now())) { // 自动结束接注 try { gameService.endBet(round.getGroupId(), 0L); // 0L表示系统自动触发 // 触发开始发红包,延迟时间从配置读取 int sendRpDelay = getConfigInt("send_rp_delay", 20); scheduleStartRp(round.getId(), round.getGroupId(), sendRpDelay); } catch (Exception e) { log.error("自动结束接注失败: roundId={}, error={}", round.getId(), e.getMessage()); } } } } /** * 每日凌晨重置活跃假人池 */ @Scheduled(cron = "0 0 0 * * ?") public void resetDailyFakeUsers() { fakeUserService.resetDailyActiveFakes(); } /** * 安排开始发红包 */ public String scheduleStartRp(Long roundId, Long groupId, int delaySeconds) { String taskId = "start_rp_" + taskIdGenerator.getAndIncrement(); PendingTask task = new PendingTask(); task.setTaskId(taskId); task.setType("start_rp"); task.setRoundId(roundId); task.setGroupId(groupId); task.setExecuteTime(LocalDateTime.now().plusSeconds(delaySeconds)); pendingTasks.put(taskId, task); log.info("安排{}秒后开始发红包: roundId={}", delaySeconds, roundId); return taskId; } /** * 安排自动结算 */ public String scheduleSettle(Long roundId, Long groupId, int digit, int delaySeconds) { String taskId = "settle_" + taskIdGenerator.getAndIncrement(); PendingTask task = new PendingTask(); task.setTaskId(taskId); task.setType("settle"); task.setRoundId(roundId); task.setGroupId(groupId); task.setDigit(digit); task.setExecuteTime(LocalDateTime.now().plusSeconds(delaySeconds)); pendingTasks.put(taskId, task); log.info("安排{}秒后结算: taskId={}, roundId={}, digit={}", delaySeconds, taskId, roundId, digit); return taskId; } /** * 安排自动下一期 */ public String scheduleAutoNextTask(Long groupId, int delaySeconds) { String taskId = "auto_next_" + taskIdGenerator.getAndIncrement(); PendingTask task = new PendingTask(); task.setTaskId(taskId); task.setType("auto_next"); task.setGroupId(groupId); task.setExecuteTime(LocalDateTime.now().plusSeconds(delaySeconds)); pendingTasks.put(taskId, task); log.info("安排{}秒后自动开始下一期: groupId={}", delaySeconds, groupId); return taskId; } /** * 安排机器人发包 */ public String scheduleBotSendRp(Long roundId, Long groupId, int delaySeconds) { String taskId = "bot_send_rp_" + taskIdGenerator.getAndIncrement(); PendingTask task = new PendingTask(); task.setTaskId(taskId); task.setType("bot_send_rp"); task.setRoundId(roundId); task.setGroupId(groupId); task.setExecuteTime(LocalDateTime.now().plusSeconds(delaySeconds)); pendingTasks.put(taskId, task); log.info("安排{}秒后机器人发包: roundId={}", delaySeconds, roundId); return taskId; } /** * 开始游戏后安排定时器 */ public void onGameStart(Long roundId, Long groupId, int intervalSeconds) { RoundTask task = new RoundTask(); task.roundId = roundId; task.groupId = groupId; task.endBetTime = LocalDateTime.now().plusSeconds(intervalSeconds); roundTasks.put(roundId, task); // 广播开始下注 GameRound round = gameRoundMapper.selectById(roundId); if (round != null) { broadcastStartBet(groupId, round.getRoundNo()); } log.info("游戏开始: roundId={}, {}秒后结束接注", roundId, intervalSeconds); } /** * 结束接注后安排发红包 */ public void onEndBet(Long roundId, Long groupId) { // 移除结束接注定时器 roundTasks.remove(roundId); // 广播结束下注 GameRound round = gameRoundMapper.selectById(roundId); if (round != null) { broadcastEndBet(groupId, round.getRoundNo()); } // 通知白名单用户尾数试算完成 notifyWhiteUsers(roundId); // 安排延迟后开始发红包 int sendRpDelay = getConfigInt("send_rp_delay", 20); scheduleStartRp(roundId, groupId, sendRpDelay); } /** * 开始发红包 */ public void onStartRp(Long roundId, Long groupId) { GameRound round = gameRoundMapper.selectById(roundId); if (round == null) { return; } // 广播开始发红包 broadcastStartRp(groupId, round.getRoundNo()); // 安排超时后如果没人发包则机器人自动发 int autoRpTimeout = getConfigInt("auto_rp_timeout", 10); scheduleBotSendRp(roundId, groupId, autoRpTimeout); } /** * 机器人自动发包 */ public void botSendRp(Long roundId, Long groupId) { GameRound round = gameRoundMapper.selectById(roundId); if (round == null || round.getTargetDigit() == null) { return; } // 查找机器人账号 User bot = userMapper.selectOne(new LambdaQueryWrapper() .eq(User::getIsBot, 1) .last("LIMIT 1")); if (bot == null) { log.error("未找到机器人账号"); return; } try { // 机器人发1元3个包 RedPacket rp = redPacketService.sendRp(bot.getId(), groupId, new BigDecimal("1"), 3, roundId); // 红包尾数控制: 需要在发红包时预计算各份额金额 // 确保手气王金额的尾数等于目标尾数 // 注: 当前实现中红包金额随机分配,实际生产需要调用RedPacketService控制尾数 log.info("机器人自动发包: roundId={}, rpId={}, targetDigit={}", roundId, rp.getId(), round != null ? round.getTargetDigit() : "null"); } catch (Exception e) { log.error("机器人发包失败: roundId={}, error={}", roundId, e.getMessage()); } } /** * 开奖结算 */ public void onSettle(Long roundId, Long groupId, int digit) { try { gameService.settle(roundId, digit); // 广播开奖结果 GameRound round = gameRoundMapper.selectById(roundId); if (round != null) { broadcastResult(groupId, digit, round.getRoundNo()); } // 检查是否自动开启下一期 GameRound roundAfterSettle = gameRoundMapper.selectById(roundId); if (roundAfterSettle != null && roundAfterSettle.getAutoNext() == 1) { // 5秒后自动开始下一期 scheduleAutoNextTask(groupId, 5); } // 清理任务 roundTasks.remove(roundId); log.info("局结算完成: roundId={}, digit={}", roundId, digit); } catch (Exception e) { log.error("结算失败: roundId={}, error={}", roundId, e.getMessage()); } } /** * 通知白名单用户(尾数试算完成) */ private void notifyWhiteUsers(Long roundId) { List whites = digitWhiteMapper.selectList(new LambdaQueryWrapper<>()); GameRound round = gameRoundMapper.selectById(roundId); if (round == null) return; for (DigitWhite white : whites) { User user = userMapper.selectById(white.getUserId()); if (user != null) { // 通过OpenIM发送通知 openIMClient.sendUserMessage(String.valueOf(user.getId()), "尾数试算完成,请前往修改尾数。期号: " + round.getRoundNo()); log.info("通知白名单用户尾数试算完成: userId={}, roundId={}", user.getId(), roundId); } } } /** * 广播开始下注 */ private void broadcastStartBet(Long groupId, String roundNo) { // 注意: 实际生产中需要通过Group表映射获取OpenIM群ID openIMClient.sendGroupMessage(String.valueOf(groupId), roundNo + "期开始下注"); } /** * 广播结束下注 */ private void broadcastEndBet(Long groupId, String roundNo) { openIMClient.sendGroupMessage(String.valueOf(groupId), roundNo + "期结束下注"); } /** * 广播开始发红包 */ private void broadcastStartRp(Long groupId, String roundNo) { openIMClient.sendGroupMessage(String.valueOf(groupId), roundNo + "期开始发红包"); } /** * 广播开奖结果 */ private void broadcastResult(Long groupId, int digit, String roundNo) { String playWinDesc = gameService.getPlayWinDesc(digit); String message = roundNo + "期开奖\n" + "尾数: " + digit + "\n" + "中奖玩法: " + playWinDesc; // 查询近30期历史 String history = getRecentHistory(groupId); if (history != null && !history.isEmpty()) { message += "\n\n近30期:\n" + history; } openIMClient.sendGroupMessage(String.valueOf(groupId), message); } /** * 获取近30期历史开奖记录 */ private String getRecentHistory(Long groupId) { List records = drawRecordMapper.selectList( new LambdaQueryWrapper() .eq(DrawRecord::getGroupId, groupId) .orderByDesc(DrawRecord::getCreateTime) .last("LIMIT 30")); if (records == null || records.isEmpty()) { return ""; } StringBuilder sb = new StringBuilder(); for (DrawRecord record : records) { sb.append(record.getRoundNo()).append(":").append(record.getDigit()).append("\n"); } return sb.toString().trim(); } /** * 局任务信息 */ private static class RoundTask { Long roundId; Long groupId; LocalDateTime endBetTime; LocalDateTime startRpTime; } /** * 延迟任务信息 */ private static class PendingTask { private String taskId; private String type; private Long roundId; private Long groupId; private Integer digit; private LocalDateTime executeTime; public String getTaskId() { return taskId; } public void setTaskId(String taskId) { this.taskId = taskId; } public String getType() { return type; } public void setType(String type) { this.type = type; } public Long getRoundId() { return roundId; } public void setRoundId(Long roundId) { this.roundId = roundId; } public Long getGroupId() { return groupId; } public void setGroupId(Long groupId) { this.groupId = groupId; } public Integer getDigit() { return digit; } public void setDigit(Integer digit) { this.digit = digit; } public LocalDateTime getExecuteTime() { return executeTime; } public void setExecuteTime(LocalDateTime executeTime) { this.executeTime = executeTime; } } /** * 处理延迟任务队列(每秒执行) */ @Scheduled(fixedDelay = 1000) public void processPendingTasks() { LocalDateTime now = LocalDateTime.now(); pendingTasks.entrySet().removeIf(entry -> { PendingTask task = entry.getValue(); if (task.getExecuteTime().isBefore(now) || task.getExecuteTime().isEqual(now)) { executeTask(task); return true; } return false; }); } /** * 执行延迟任务 */ private void executeTask(PendingTask task) { try { switch (task.getType()) { case "settle": onSettle(task.getRoundId(), task.getGroupId(), task.getDigit()); break; case "auto_next": onAutoNext(task.getGroupId()); break; case "start_rp": onStartRp(task.getRoundId(), task.getGroupId()); break; case "bot_send_rp": botSendRp(task.getRoundId(), task.getGroupId()); break; default: log.warn("未知任务类型: {}", task.getType()); } } catch (Exception e) { log.error("执行任务失败: taskId={}, error={}", task.getTaskId(), e.getMessage()); } } /** * 自动开始下一期 */ private void onAutoNext(Long groupId) { try { gameService.startGame(groupId, 0L); // 0L表示系统自动 } catch (Exception e) { log.error("自动开始下一期失败: groupId={}, error={}", groupId, e.getMessage()); } } }