MCPcopy Create free account
hub / github.com/e2wugui/zeze / fireCron

Method fireCron

ZezeJava/ZezeJava/src/main/java/Zeze/Component/Timer.java:1099–1166  ·  view source on GitHub ↗
(long timerSerialId, int serverId, @NotNull String timerId, long concurrentSerialNo,
						  boolean missfire)

Source from the content-addressed store, hash-verified

1097 }
1098
1099 private void fireCron(long timerSerialId, int serverId, @NotNull String timerId, long concurrentSerialNo,
1100 boolean missfire) {
1101 if (Task.call(zeze.newProcedure(() -> {
1102 var index = _tIndexs.get(timerId);
1103 if (index == null
1104 || index.getServerId() != zeze.getConfig().getServerId() // 不是拥有者,取消本地调度,应该是不大可能发生的。
1105 || index.getSerialId() != timerSerialId // 新注册的,旧的future需要取消。
1106 ) {
1107 Transaction.whileCommit(() -> cancelFuture(timerId));
1108 return 0;
1109 }
1110
1111 var nodeId = index.getNodeId();
1112 var node = _tNodes.get(nodeId);
1113 if (node == null) {
1114 cancel(serverId, timerId, nodeId, null, null); // maybe concurrent cancel
1115 return 0;
1116 }
1117
1118 var timer = node.getTimers().get(timerId);
1119 if (timer == null)
1120 throw new IllegalStateException("maybe operate before timer created");
1121 var handle = findTimerHandle(timer.getHandleName());
1122 var cronTimer = timer.getTimerObj_Zeze_Builtin_Timer_BCronTimer();
1123 if (concurrentSerialNo == timer.getConcurrentFireSerialNo()) {
1124 var hasNext = Timer.nextCronTimer(cronTimer, missfire);
1125 var context = new TimerContext(this, timer, cronTimer.getHappenTimes(),
1126 cronTimer.getExpectedTime(), cronTimer.getNextExpectedTime());
1127
1128 // 当调度发生了错误或者由于异步时序没有原子保证,导致同时(或某个瞬间)在多个Server进程调度时,
1129 // 这个系列号保证触发用户回调只会发生一次。这个并发问题不取消定时器,继续尝试调度(去争抢执行权)。
1130 // 定时器的调度生命期由其他地方保证最终一致。如果保证发生了错误,将一致并发争抢执行权。
1131 var serialSaved = index.getSerialId();
1132 var ret = Task.call(zeze.newProcedure(() -> {
1133 handle.onTimer(context);
1134 return 0;
1135 }, "Timer.fireCronUser." + timer.getHandleName()));
1136
1137 var indexNew = _tIndexs.get(timerId);
1138 if (indexNew == null || indexNew.getSerialId() != serialSaved)
1139 return 0; // canceled or new timer.
1140
1141 if (ret == Procedure.Exception) {
1142 // 用户处理不允许异常,其他错误记录忽略,日志已经记录。
1143 cancel(serverId, timerId, nodeId, node, handle);
1144 return 0;
1145 }
1146 timer.setConcurrentFireSerialNo(concurrentSerialNo + 1);
1147
1148 if (!hasNext) {
1149 cancel(serverId, timerId, nodeId, node, handle);
1150 return 0;
1151 }
1152 }
1153 // else 发生了并发执行争抢,也需要再次进行本地调度。此时直接使用cronTimer中的值,不需要再次进行计算。
1154
1155 // continue period
1156 scheduleCronNext(timerSerialId, serverId, timerId,

Callers 2

scheduleCronNextMethod · 0.95
loadTimerMethod · 0.95

Calls 15

callMethod · 0.95
whileCommitMethod · 0.95
cancelFutureMethod · 0.95
cancelMethod · 0.95
findTimerHandleMethod · 0.95
nextCronTimerMethod · 0.95
scheduleCronNextMethod · 0.95
getTimersMethod · 0.80
getMethod · 0.65
getServerIdMethod · 0.65
getSerialIdMethod · 0.65

Tested by

no test coverage detected