1use super::{
2 Arc, AtomicI64, AtomicOrdering, BTreeSet, BinaryHeap, BlockPos, BlockRef, ChunkPos,
3 ChunkTickContainer, ChunkTickContainerLifecycle, ChunkTickLists, FluidRef, FullChunkRef,
4 FxHashMap, Ordering, PackedChunkPos, ScheduledTick, ScheduledTickBatch, SyncMutex, SyncRwLock,
5 TickKey, TickList, TickPriority, TickSchedulerError, intra_tick_drain_order,
6};
7
8pub(crate) struct WorldTickScheduler {
9 next_sub_tick_order: AtomicI64,
10 phase: SyncRwLock<()>,
11 pub(super) state: SyncMutex<WorldTickSchedulerState>,
12}
13
14#[derive(Debug, Default)]
15pub(super) struct WorldTickSchedulerState {
16 pub(super) chunks: FxHashMap<ChunkPos, RegisteredChunkTicks>,
17 pub(super) active_block_deadlines: BTreeSet<(i64, PackedChunkPos)>,
18 active_fluid_deadlines: BTreeSet<(i64, PackedChunkPos)>,
19 pub(super) active_generation: u64,
20}
21
22#[derive(Debug)]
23pub(super) struct RegisteredChunkTicks {
24 pub(super) container: Arc<ChunkTickContainer>,
25 pub(super) block_head: Option<i64>,
26 pub(super) fluid_head: Option<i64>,
27 pub(super) active: Option<ActiveChunkRank>,
28}
29
30#[derive(Debug, Clone, Copy)]
31pub(super) struct ActiveChunkRank {
32 generation: u64,
33 rank: usize,
34}
35
36#[derive(Clone, Copy)]
37pub(super) enum TickKind {
38 Block,
39 Fluid,
40}
41
42struct ActiveTickContainer {
43 pos: ChunkPos,
44 rank: usize,
45 container: Arc<ChunkTickContainer>,
46}
47
48struct TickHeadUpdate {
49 pos: ChunkPos,
50 trigger_tick: Option<i64>,
51}
52
53impl RegisteredChunkTicks {
54 const fn head(&self, kind: TickKind) -> Option<i64> {
55 match kind {
56 TickKind::Block => self.block_head,
57 TickKind::Fluid => self.fluid_head,
58 }
59 }
60
61 const fn set_head(&mut self, kind: TickKind, trigger_tick: Option<i64>) {
62 match kind {
63 TickKind::Block => self.block_head = trigger_tick,
64 TickKind::Fluid => self.fluid_head = trigger_tick,
65 }
66 }
67
68 pub(super) const fn active_rank(&self, generation: u64) -> Option<usize> {
69 match self.active {
70 Some(active) if active.generation == generation => Some(active.rank),
71 _ => None,
72 }
73 }
74}
75
76impl WorldTickSchedulerState {
77 fn advance_active_generation(&mut self) -> u64 {
78 let generation = if let Some(generation) = self.active_generation.checked_add(1) {
79 generation
80 } else {
81 for registered in self.chunks.values_mut() {
82 registered.active = None;
83 }
84 1
85 };
86 self.active_generation = generation;
87 generation
88 }
89
90 const fn deadlines(&self, kind: TickKind) -> &BTreeSet<(i64, PackedChunkPos)> {
91 match kind {
92 TickKind::Block => &self.active_block_deadlines,
93 TickKind::Fluid => &self.active_fluid_deadlines,
94 }
95 }
96
97 const fn deadlines_mut(&mut self, kind: TickKind) -> &mut BTreeSet<(i64, PackedChunkPos)> {
98 match kind {
99 TickKind::Block => &mut self.active_block_deadlines,
100 TickKind::Fluid => &mut self.active_fluid_deadlines,
101 }
102 }
103
104 pub(super) fn set_head(
105 &mut self,
106 pos: ChunkPos,
107 kind: TickKind,
108 trigger_tick: Option<i64>,
109 ) -> Result<(), TickSchedulerError> {
110 let Some(registered) = self.chunks.get(&pos) else {
111 return Err(TickSchedulerError::MissingContainer(pos));
112 };
113 let active = registered.active_rank(self.active_generation).is_some();
114 let previous = registered.head(kind);
115 if previous == trigger_tick {
116 return Ok(());
117 }
118 let packed = PackedChunkPos::from(pos);
119 if active && let Some(previous) = previous {
120 assert!(
121 self.deadlines_mut(kind).remove(&(previous, packed)),
122 "active scheduled-tick head was absent from its deadline index"
123 );
124 }
125 let Some(registered) = self.chunks.get_mut(&pos) else {
126 return Err(TickSchedulerError::MissingContainer(pos));
127 };
128 registered.set_head(kind, trigger_tick);
129 if active && let Some(trigger_tick) = trigger_tick {
130 assert!(
131 self.deadlines_mut(kind).insert((trigger_tick, packed)),
132 "active scheduled-tick deadline was already indexed"
133 );
134 }
135 Ok(())
136 }
137
138 fn take_due(&mut self, kind: TickKind, current_tick: i64) -> Vec<ActiveTickContainer> {
139 let mut due = Vec::new();
140 while let Some((trigger_tick, packed)) = self.deadlines(kind).first().copied() {
141 if trigger_tick > current_tick {
142 break;
143 }
144 assert_eq!(
145 self.deadlines_mut(kind).pop_first(),
146 Some((trigger_tick, packed)),
147 "due scheduled-tick deadline changed during collection"
148 );
149 let pos = packed.to_chunk_pos();
150 let Some(registered) = self.chunks.get_mut(&pos) else {
151 panic!("active scheduled-tick deadline lost its registered container");
152 };
153 assert_eq!(
154 registered.head(kind),
155 Some(trigger_tick),
156 "active scheduled-tick deadline diverged from its container head"
157 );
158 let Some(rank) = registered.active_rank(self.active_generation) else {
159 panic!("inactive scheduled-tick container retained an active deadline");
160 };
161 due.push(ActiveTickContainer {
162 pos,
163 rank,
164 container: Arc::clone(®istered.container),
165 });
166 registered.set_head(kind, None);
167 }
168 due
169 }
170}
171
172impl WorldTickScheduler {
173 #[must_use]
174 pub(crate) fn new() -> Self {
175 Self {
176 next_sub_tick_order: AtomicI64::new(0),
177 phase: SyncRwLock::new(()),
178 state: SyncMutex::new(WorldTickSchedulerState::default()),
179 }
180 }
181
182 pub(crate) fn next_sub_tick_order(&self) -> i64 {
187 self.next_sub_tick_order
188 .fetch_add(1, AtomicOrdering::Relaxed)
189 }
190
191 pub(crate) fn register_chunk(&self, chunk: FullChunkRef<'_>) -> Result<(), TickSchedulerError> {
193 let _phase = self.phase.read();
194 let container = chunk.scheduled_tick_container();
195 let mut container_state = container.state.lock();
196 if container_state.lifecycle != ChunkTickContainerLifecycle::PrePublication {
197 return Err(TickSchedulerError::AlreadyRegistered(chunk.common().pos));
198 }
199 let block_head = container_state
200 .lists
201 .block()
202 .peek()
203 .map(|tick| tick.trigger_tick);
204 let fluid_head = container_state
205 .lists
206 .fluid()
207 .peek()
208 .map(|tick| tick.trigger_tick);
209 let mut state = self.state.lock();
210 if state.chunks.contains_key(&chunk.common().pos) {
211 return Err(TickSchedulerError::AlreadyRegistered(chunk.common().pos));
212 }
213 state.chunks.insert(
214 chunk.common().pos,
215 RegisteredChunkTicks {
216 container: Arc::clone(container),
217 block_head,
218 fluid_head,
219 active: None,
220 },
221 );
222 container_state.lifecycle = ChunkTickContainerLifecycle::Registered;
223 Ok(())
224 }
225
226 pub(crate) fn unpack_chunk(
231 &self,
232 pos: ChunkPos,
233 current_tick: i64,
234 ) -> Result<(), TickSchedulerError> {
235 let _phase = self.phase.read();
236 let container = self
237 .state
238 .lock()
239 .chunks
240 .get(&pos)
241 .map(|registered| Arc::clone(®istered.container))
242 .ok_or(TickSchedulerError::MissingContainer(pos))?;
243 let mut container_state = container.state.lock();
244 if container_state.lifecycle != ChunkTickContainerLifecycle::Registered {
245 return Err(TickSchedulerError::MissingContainer(pos));
246 }
247 container_state.lists.block_mut().unpack(current_tick);
248 container_state.lists.fluid_mut().unpack(current_tick);
249 let block_head = container_state
250 .lists
251 .block()
252 .peek()
253 .map(|tick| tick.trigger_tick);
254 let fluid_head = container_state
255 .lists
256 .fluid()
257 .peek()
258 .map(|tick| tick.trigger_tick);
259 let mut state = self.state.lock();
260 let Some(registered) = state.chunks.get(&pos) else {
261 return Err(TickSchedulerError::MissingContainer(pos));
262 };
263 if !Arc::ptr_eq(®istered.container, &container) {
264 return Err(TickSchedulerError::ContainerMismatch(pos));
265 }
266 state.set_head(pos, TickKind::Block, block_head)?;
267 state.set_head(pos, TickKind::Fluid, fluid_head)?;
268 Ok(())
269 }
270
271 pub(crate) fn unregister_chunk(&self, pos: ChunkPos) {
273 let _phase = self.phase.write();
274 let registered = {
275 let mut state = self.state.lock();
276 let active_generation = state.active_generation;
277 let registered = state.chunks.remove(&pos);
278 if let Some(registered) = ®istered
279 && registered.active_rank(active_generation).is_some()
280 {
281 let packed = PackedChunkPos::from(pos);
282 if let Some(head) = registered.block_head {
283 assert!(
284 state.active_block_deadlines.remove(&(head, packed)),
285 "unloaded active block-tick head was absent from its deadline index"
286 );
287 }
288 if let Some(head) = registered.fluid_head {
289 assert!(
290 state.active_fluid_deadlines.remove(&(head, packed)),
291 "unloaded active fluid-tick head was absent from its deadline index"
292 );
293 }
294 }
295 registered
296 };
297 if let Some(registered) = registered {
298 registered.container.state.lock().lifecycle = ChunkTickContainerLifecycle::Finalized;
299 }
300 }
301
302 #[cfg(test)]
303 pub(crate) fn has_registered_chunk(&self, pos: ChunkPos) -> bool {
304 self.state.lock().chunks.contains_key(&pos)
305 }
306
307 #[cfg(test)]
308 pub(crate) fn has_indexed_head(&self, pos: ChunkPos) -> bool {
309 let state = self.state.lock();
310 state.chunks.get(&pos).is_some_and(|registered| {
311 registered.block_head.is_some() || registered.fluid_head.is_some()
312 })
313 }
314
315 pub(crate) fn reconcile_active_chunks<I>(
317 &self,
318 active_chunks: I,
319 ) -> Result<(), TickSchedulerError>
320 where
321 I: Iterator<Item = ChunkPos> + Clone,
322 {
323 let _phase = self.phase.write();
324 let mut state = self.state.lock();
325 for pos in active_chunks.clone() {
326 if !state.chunks.contains_key(&pos) {
327 return Err(TickSchedulerError::MissingContainer(pos));
328 }
329 }
330 state.active_block_deadlines.clear();
331 state.active_fluid_deadlines.clear();
332 let generation = state.advance_active_generation();
333 for (rank, pos) in active_chunks.enumerate() {
334 let (block_head, fluid_head) = {
335 let Some(registered) = state.chunks.get_mut(&pos) else {
336 return Err(TickSchedulerError::MissingContainer(pos));
337 };
338 assert!(
339 registered
340 .active
341 .is_none_or(|active| active.generation != generation),
342 "active scheduled-tick chunk appeared twice during reconciliation"
343 );
344 registered.active = Some(ActiveChunkRank { generation, rank });
345 (registered.block_head, registered.fluid_head)
346 };
347 let packed = PackedChunkPos::from(pos);
348 if let Some(block_head) = block_head {
349 assert!(
350 state.active_block_deadlines.insert((block_head, packed)),
351 "active block-tick chunk appeared twice during reconciliation"
352 );
353 }
354 if let Some(fluid_head) = fluid_head {
355 assert!(
356 state.active_fluid_deadlines.insert((fluid_head, packed)),
357 "active fluid-tick chunk appeared twice during reconciliation"
358 );
359 }
360 }
361 Ok(())
362 }
363
364 pub(crate) fn schedule_block(
365 &self,
366 chunk: FullChunkRef<'_>,
367 block: BlockRef,
368 pos: BlockPos,
369 trigger_tick: i64,
370 priority: TickPriority,
371 sub_tick_order: i64,
372 ) -> Result<bool, TickSchedulerError> {
373 let _phase = self.phase.read();
374 let container = chunk.scheduled_tick_container();
375 let mut container_state = container.state.lock();
376 if container_state.lifecycle == ChunkTickContainerLifecycle::PrePublication {
377 return Ok(container_state.lists.block_mut().schedule(
378 block,
379 pos,
380 trigger_tick,
381 priority,
382 sub_tick_order,
383 ));
384 }
385 if container_state.lifecycle == ChunkTickContainerLifecycle::Finalized {
386 return Err(TickSchedulerError::MissingContainer(chunk.common().pos));
387 }
388 let previous_head = container_state
389 .lists
390 .block()
391 .peek()
392 .map(|tick| tick.trigger_tick);
393 let added = container_state.lists.block_mut().schedule(
394 block,
395 pos,
396 trigger_tick,
397 priority,
398 sub_tick_order,
399 );
400 if !added {
401 return Ok(false);
402 }
403 let head = container_state
404 .lists
405 .block()
406 .peek()
407 .map(|tick| tick.trigger_tick);
408 if head == previous_head {
409 return Ok(true);
410 }
411 let mut state = self.state.lock();
412 let Some(registered) = state.chunks.get(&chunk.common().pos) else {
413 return Err(TickSchedulerError::MissingContainer(chunk.common().pos));
414 };
415 if !Arc::ptr_eq(®istered.container, container) {
416 return Err(TickSchedulerError::ContainerMismatch(chunk.common().pos));
417 }
418 state.set_head(chunk.common().pos, TickKind::Block, head)?;
419 Ok(true)
420 }
421
422 pub(crate) fn schedule_fluid(
423 &self,
424 chunk: FullChunkRef<'_>,
425 fluid: FluidRef,
426 pos: BlockPos,
427 trigger_tick: i64,
428 priority: TickPriority,
429 sub_tick_order: i64,
430 ) -> Result<bool, TickSchedulerError> {
431 let _phase = self.phase.read();
432 let container = chunk.scheduled_tick_container();
433 let mut container_state = container.state.lock();
434 if container_state.lifecycle == ChunkTickContainerLifecycle::PrePublication {
435 return Ok(container_state.lists.fluid_mut().schedule(
436 fluid,
437 pos,
438 trigger_tick,
439 priority,
440 sub_tick_order,
441 ));
442 }
443 if container_state.lifecycle == ChunkTickContainerLifecycle::Finalized {
444 return Err(TickSchedulerError::MissingContainer(chunk.common().pos));
445 }
446 let previous_head = container_state
447 .lists
448 .fluid()
449 .peek()
450 .map(|tick| tick.trigger_tick);
451 let added = container_state.lists.fluid_mut().schedule(
452 fluid,
453 pos,
454 trigger_tick,
455 priority,
456 sub_tick_order,
457 );
458 if !added {
459 return Ok(false);
460 }
461 let head = container_state
462 .lists
463 .fluid()
464 .peek()
465 .map(|tick| tick.trigger_tick);
466 if head == previous_head {
467 return Ok(true);
468 }
469 let mut state = self.state.lock();
470 let Some(registered) = state.chunks.get(&chunk.common().pos) else {
471 return Err(TickSchedulerError::MissingContainer(chunk.common().pos));
472 };
473 if !Arc::ptr_eq(®istered.container, container) {
474 return Err(TickSchedulerError::ContainerMismatch(chunk.common().pos));
475 }
476 state.set_head(chunk.common().pos, TickKind::Fluid, head)?;
477 Ok(true)
478 }
479
480 pub(crate) fn begin_tick(
482 &self,
483 current_tick: i64,
484 max_ticks: usize,
485 ) -> ScheduledTickBatch<BlockRef> {
486 self.collect_ticks(
487 TickKind::Block,
488 ChunkTickLists::block_mut,
489 current_tick,
490 max_ticks,
491 )
492 }
493
494 pub(crate) fn collect_fluid_ticks(
496 &self,
497 current_tick: i64,
498 max_ticks: usize,
499 ) -> ScheduledTickBatch<FluidRef> {
500 self.collect_ticks(
501 TickKind::Fluid,
502 ChunkTickLists::fluid_mut,
503 current_tick,
504 max_ticks,
505 )
506 }
507
508 fn collect_ticks<T: TickKey>(
509 &self,
510 kind: TickKind,
511 select: fn(&mut ChunkTickLists) -> &mut TickList<T>,
512 current_tick: i64,
513 max_ticks: usize,
514 ) -> ScheduledTickBatch<T> {
515 if max_ticks == 0 {
516 return ScheduledTickBatch {
517 ticks: Vec::new(),
518 changed_containers: Vec::new(),
519 };
520 }
521 let _phase = self.phase.write();
522 let due = self.state.lock().take_due(kind, current_tick);
523 if due.is_empty() {
524 return ScheduledTickBatch {
525 ticks: Vec::new(),
526 changed_containers: Vec::new(),
527 };
528 }
529 let (batch, head_updates) = collect_registered_ticks(due, select, current_tick, max_ticks);
530 let mut state = self.state.lock();
531 for update in head_updates {
532 if let Err(error) = state.set_head(update.pos, kind, update.trigger_tick) {
533 panic!("scheduled-tick head index invariant failed: {error:?}");
534 }
535 }
536 batch
537 }
538}
539
540impl Default for WorldTickScheduler {
541 fn default() -> Self {
542 Self::new()
543 }
544}
545
546#[derive(Debug)]
547struct ReadyContainer {
548 pos: ChunkPos,
549 rank: usize,
550 container: Arc<ChunkTickContainer>,
551 priority: TickPriority,
552 sub_tick_order: i64,
553 dirty_reported: bool,
554}
555
556impl PartialEq for ReadyContainer {
557 fn eq(&self, other: &Self) -> bool {
558 self.priority == other.priority
559 && self.sub_tick_order == other.sub_tick_order
560 && self.rank == other.rank
561 }
562}
563
564impl Eq for ReadyContainer {}
565
566impl ReadyContainer {
567 fn new<T: TickKey>(
568 active: ActiveTickContainer,
569 tick: ScheduledTick<T>,
570 dirty_reported: bool,
571 ) -> Self {
572 Self {
573 pos: active.pos,
574 rank: active.rank,
575 container: active.container,
576 priority: tick.priority,
577 sub_tick_order: tick.sub_tick_order,
578 dirty_reported,
579 }
580 }
581
582 const fn refresh<T: TickKey>(&mut self, tick: ScheduledTick<T>) {
583 self.priority = tick.priority;
584 self.sub_tick_order = tick.sub_tick_order;
585 }
586}
587
588impl PartialOrd for ReadyContainer {
589 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
590 Some(self.cmp(other))
591 }
592}
593
594impl Ord for ReadyContainer {
595 fn cmp(&self, other: &Self) -> Ordering {
596 intra_tick_drain_order(
597 self.priority,
598 self.sub_tick_order,
599 other.priority,
600 other.sub_tick_order,
601 )
602 .reverse()
603 .then_with(|| other.rank.cmp(&self.rank))
604 }
605}
606
607fn prepare_ready_containers<T: TickKey>(
608 due_containers: Vec<ActiveTickContainer>,
609 select: fn(&mut ChunkTickLists) -> &mut TickList<T>,
610 current_tick: i64,
611) -> (BinaryHeap<ReadyContainer>, Vec<TickHeadUpdate>) {
612 let mut ready_containers = BinaryHeap::with_capacity(due_containers.len());
613 let mut head_updates = Vec::with_capacity(due_containers.len());
614 for active in due_containers {
615 let head = {
616 let mut state = active.container.state.lock();
617 assert_eq!(
618 state.lifecycle,
619 ChunkTickContainerLifecycle::Registered,
620 "active scheduled-tick container was not registered"
621 );
622 select(&mut state.lists).peek()
623 };
624 if let Some(tick) = head
625 && tick.trigger_tick <= current_tick
626 {
627 ready_containers.push(ReadyContainer::new(active, tick, false));
628 } else {
629 head_updates.push(TickHeadUpdate {
630 pos: active.pos,
631 trigger_tick: head.map(|tick| tick.trigger_tick),
632 });
633 }
634 }
635 (ready_containers, head_updates)
636}
637
638fn collect_registered_ticks<T: TickKey>(
644 due_containers: Vec<ActiveTickContainer>,
645 select: fn(&mut ChunkTickLists) -> &mut TickList<T>,
646 current_tick: i64,
647 max_ticks: usize,
648) -> (ScheduledTickBatch<T>, Vec<TickHeadUpdate>) {
649 let (mut ready_containers, mut head_updates) =
650 prepare_ready_containers(due_containers, select, current_tick);
651
652 let mut ticks = Vec::with_capacity(max_ticks.min(ready_containers.len()));
653 let mut changed_containers = Vec::with_capacity(ready_containers.len());
654 while ticks.len() < max_ticks {
655 let Some(mut ready_container) = ready_containers.pop() else {
656 break;
657 };
658 let mut container_state = ready_container.container.state.lock();
659 assert_eq!(
660 container_state.lifecycle,
661 ChunkTickContainerLifecycle::Registered,
662 "ready scheduled-tick container was not registered"
663 );
664 let container = select(&mut container_state.lists);
665 let Some(tick) = container.pop_ready(current_tick) else {
666 head_updates.push(TickHeadUpdate {
667 pos: ready_container.pos,
668 trigger_tick: container.peek().map(|tick| tick.trigger_tick),
669 });
670 continue;
671 };
672
673 if !ready_container.dirty_reported {
674 changed_containers.push(ready_container.rank);
675 ready_container.dirty_reported = true;
676 }
677 ticks.push(tick);
678
679 let next_competing_container = ready_containers
683 .peek()
684 .map(|competitor| (competitor.priority, competitor.sub_tick_order));
685 while ticks.len() < max_ticks {
686 let Some(next_tick) = container.peek_ready(current_tick) else {
687 break;
688 };
689 if next_competing_container.is_some_and(|(priority, sub_tick_order)| {
690 intra_tick_drain_order(
691 next_tick.priority,
692 next_tick.sub_tick_order,
693 priority,
694 sub_tick_order,
695 ) == Ordering::Greater
696 }) {
697 break;
698 }
699 let Some(next_tick) = container.pop_ready(current_tick) else {
700 break;
701 };
702 ticks.push(next_tick);
703 }
704
705 let next_tick = container.peek();
706 drop(container_state);
707 if let Some(next_tick) = next_tick {
708 if ticks.len() < max_ticks && next_tick.trigger_tick <= current_tick {
709 ready_container.refresh(next_tick);
710 ready_containers.push(ready_container);
711 } else {
712 head_updates.push(TickHeadUpdate {
713 pos: ready_container.pos,
714 trigger_tick: Some(next_tick.trigger_tick),
715 });
716 }
717 } else {
718 head_updates.push(TickHeadUpdate {
719 pos: ready_container.pos,
720 trigger_tick: None,
721 });
722 }
723 }
724
725 for ready_container in ready_containers {
726 let next_tick = {
727 let mut state = ready_container.container.state.lock();
728 select(&mut state.lists).peek()
729 };
730 head_updates.push(TickHeadUpdate {
731 pos: ready_container.pos,
732 trigger_tick: next_tick.map(|tick| tick.trigger_tick),
733 });
734 }
735
736 (
737 ScheduledTickBatch {
738 ticks,
739 changed_containers,
740 },
741 head_updates,
742 )
743}