(long timerSerialId, int serverId, @NotNull String timerId, long concurrentSerialNo, boolean missfire)
| 956 | } |
| 957 | |
| 958 | private void fireSimple(long timerSerialId, int serverId, @NotNull String timerId, long concurrentSerialNo, |
| 959 | boolean missfire) { |
| 960 | if (Task.call(zeze.newProcedure(() -> { |
| 961 | var index = _tIndexs.get(timerId); |
| 962 | if (index == null |
| 963 | || index.getServerId() != zeze.getConfig().getServerId() // 不是拥有者,取消本地调度,应该是不大可能发生的。 |
| 964 | || index.getSerialId() != timerSerialId // 新注册的,旧的future需要取消。 |
| 965 | ) { |
| 966 | Transaction.whileCommit(() -> cancelFuture(timerId)); |
| 967 | return 0; |
| 968 | } |
| 969 | |
| 970 | var nodeId = index.getNodeId(); |
| 971 | var node = _tNodes.get(nodeId); |
| 972 | if (node == null) { |
| 973 | cancel(serverId, timerId, nodeId, null, null); |
| 974 | return 0; |
| 975 | } |
| 976 | |
| 977 | var timer = node.getTimers().get(timerId); |
| 978 | if (timer == null) |
| 979 | throw new IllegalStateException("maybe operate before timer created"); |
| 980 | var handle = findTimerHandle(timer.getHandleName()); |
| 981 | var simpleTimer = timer.getTimerObj_Zeze_Builtin_Timer_BSimpleTimer(); |
| 982 | if (concurrentSerialNo == timer.getConcurrentFireSerialNo()) { |
| 983 | var hasNext = nextSimpleTimer(simpleTimer, missfire); |
| 984 | var context = new TimerContext(this, timer, simpleTimer.getHappenTimes(), |
| 985 | simpleTimer.getExpectedTime(), simpleTimer.getNextExpectedTime()); |
| 986 | |
| 987 | // 当调度发生了错误或者由于异步时序没有原子保证,导致同时(或某个瞬间)在多个Server进程调度时, |
| 988 | // 这个系列号保证触发用户回调只会发生一次。这个并发问题不取消定时器,继续尝试调度(去争抢执行权)。 |
| 989 | // 定时器的调度生命期由其他地方保证最终一致。如果保证发生了错误,将一致并发争抢执行权。 |
| 990 | var serialSaved = index.getSerialId(); |
| 991 | var ret = Task.call(zeze.newProcedure(() -> { |
| 992 | handle.onTimer(context); |
| 993 | return 0; |
| 994 | }, "Timer.fireSimpleUser." + timer.getHandleName())); |
| 995 | |
| 996 | var indexNew = _tIndexs.get(timerId); |
| 997 | if (indexNew == null || indexNew.getSerialId() != serialSaved) |
| 998 | return 0; // canceled or new timer. |
| 999 | |
| 1000 | if (ret == Procedure.Exception) { |
| 1001 | // 用户处理不允许异常,其他错误记录忽略,日志已经记录。 |
| 1002 | cancel(serverId, timerId, nodeId, node, handle); |
| 1003 | return 0; |
| 1004 | } |
| 1005 | timer.setConcurrentFireSerialNo(concurrentSerialNo + 1); |
| 1006 | // 其他错误忽略 |
| 1007 | if (!hasNext) { |
| 1008 | cancel(serverId, timerId, nodeId, node, handle); |
| 1009 | return 0; |
| 1010 | } |
| 1011 | } |
| 1012 | // else 发生了并发执行争抢,也需要再次进行本地调度。此时直接使用simpleTimer中的值,不需要再次进行计算。 |
| 1013 | |
| 1014 | // continue period |
| 1015 | scheduleSimple(timerSerialId, serverId, timerId, |
no test coverage detected