This commit is contained in:
wells
2026-08-05 15:43:04 +08:00
parent 50005dd439
commit ce7a619791
8 changed files with 283 additions and 35 deletions
+153 -26
View File
@@ -44,6 +44,12 @@ public class GameJob {
// 局状态监控(使用内存缓存,实际生产应使用Redis)
private final Map<Long, RoundTask> roundTasks = new ConcurrentHashMap<>();
// 延迟任务队列
private final Map<String, PendingTask> pendingTasks = new ConcurrentHashMap<>();
// 任务ID生成器
private java.util.concurrent.atomic.AtomicLong taskIdGenerator = new java.util.concurrent.atomic.AtomicLong(1);
/**
* 检查红包是否需要托领取(每3秒执行)
*/
@@ -131,20 +137,65 @@ public class GameJob {
/**
* 安排开始发红包
*/
public void scheduleStartRp(Long roundId, Long groupId, int delaySeconds) {
RoundTask task = new RoundTask();
task.roundId = roundId;
task.groupId = groupId;
task.startRpTime = LocalDateTime.now().plusSeconds(delaySeconds);
roundTasks.put(roundId, task);
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 void scheduleSettle(Long roundId, Long groupId, int digit, int delaySeconds) {
// TODO: 使用延迟队列或定时任务实现
log.info("安排{}秒后结算 roundId={}, digit={}", delaySeconds, roundId, digit);
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;
}
/**
@@ -199,7 +250,7 @@ public class GameJob {
broadcastStartRp(groupId, round.getRoundNo());
// 安排10秒后如果没人发包则机器人自动发
// TODO: 实现延迟检查逻辑
scheduleBotSendRp(roundId, groupId, 10);
}
/**
@@ -226,10 +277,11 @@ public class GameJob {
RedPacket rp = redPacketService.sendRp(bot.getId(), groupId,
new BigDecimal("1"), 3, roundId);
// 设置目标尾数
// TODO: 实现红包尾数控制
log.info("机器人自动发包: roundId={}, rpId={}", roundId, rp.getId());
// 红包尾数控制: 需要在发红包时预计算各份额金额
// 确保手气王金额的尾数等于目标尾数
// 注: 当前实现中红包金额随机分配,实际生产需要调用RedPacketService控制尾数
log.info("机器人自动发包: roundId={}, rpId={}, targetDigit={}",
roundId, rp.getId(), round != null ? round.getTargetDigit() : "null");
} catch (Exception e) {
log.error("机器人发包失败: roundId={}, error={}", roundId, e.getMessage());
}
@@ -252,7 +304,7 @@ public class GameJob {
GameRound roundAfterSettle = gameRoundMapper.selectById(roundId);
if (roundAfterSettle != null && roundAfterSettle.getAutoNext() == 1) {
// 5秒后自动开始下一期
scheduleAutoNext(groupId, 5);
scheduleAutoNextTask(groupId, 5);
}
// 清理任务
@@ -269,28 +321,25 @@ public class GameJob {
*/
private void notifyWhiteUsers(Long roundId) {
List<DigitWhite> whites = digitWhiteMapper.selectList(new LambdaQueryWrapper<DigitWhite>());
GameRound round = gameRoundMapper.selectById(roundId);
if (round == null) return;
for (DigitWhite white : whites) {
User user = userMapper.selectById(white.getUserId());
if (user != null) {
// TODO: 发送推送通知
log.info("通知白名单用户尾数试算完成: userId={}", user.getId());
// 通过OpenIM发送通知
openIMClient.sendUserMessage(String.valueOf(user.getId()),
"尾数试算完成,请前往修改尾数。期号: " + round.getRoundNo());
log.info("通知白名单用户尾数试算完成: userId={}, roundId={}", user.getId(), roundId);
}
}
}
/**
* 安排自动开始下一期
*/
private void scheduleAutoNext(Long groupId, int delaySeconds) {
// TODO: 实现延迟自动开局
log.info("安排{}秒后自动开始下一期", delaySeconds);
}
/**
* 广播开始下注
*/
private void broadcastStartBet(Long groupId, String roundNo) {
// TODO: 获取群对应的OpenIM群ID并发送
// 注意: 实际生产中需要通过Group表映射获取OpenIM群ID
openIMClient.sendGroupMessage(String.valueOf(groupId),
roundNo + "期开始下注");
}
@@ -331,4 +380,82 @@ public class GameJob {
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());
}
}
}