1use std::{
4 mem,
5 sync::{
6 Arc, Weak,
7 atomic::{AtomicBool, Ordering},
8 },
9};
10
11use arc_swap::ArcSwapOption;
12use rustc_hash::FxHashMap;
13use steel_registry::blocks::block_state_ext::BlockStateExt as _;
14use steel_utils::{BlockPos, locks::SyncMutex};
15
16use crate::{
17 block_entity::{BlockEntityLifecycleExt as _, BlockEntityTicker, SharedBlockEntity},
18 chunk::{chunk_holder::ChunkHolder, chunk_ticket_manager::ChunkTicketLevel},
19};
20
21use super::{World, border::WorldBorderSnapshot};
22
23pub(crate) struct BoundBlockEntityTicker {
25 holder: Weak<ChunkHolder>,
26 entity: SharedBlockEntity,
27 ticker: BlockEntityTicker,
28 logged_invalid_state: AtomicBool,
29}
30
31impl BoundBlockEntityTicker {
32 fn new(
33 holder: &Arc<ChunkHolder>,
34 entity: SharedBlockEntity,
35 ticker: BlockEntityTicker,
36 ) -> Self {
37 Self {
38 holder: Arc::downgrade(holder),
39 entity,
40 ticker,
41 logged_invalid_state: AtomicBool::new(false),
42 }
43 }
44
45 fn belongs_to(&self, holder: &Arc<ChunkHolder>) -> bool {
46 self.holder.as_ptr() == Arc::as_ptr(holder)
47 }
48}
49
50pub(crate) struct RebindableBlockEntityTicker {
52 binding: ArcSwapOption<BoundBlockEntityTicker>,
53}
54
55impl RebindableBlockEntityTicker {
56 fn new(binding: Arc<BoundBlockEntityTicker>) -> Self {
57 Self {
58 binding: ArcSwapOption::from(Some(binding)),
59 }
60 }
61
62 fn rebind(&self, binding: Option<Arc<BoundBlockEntityTicker>>) {
63 self.binding.store(binding);
64 }
65}
66
67#[derive(Default)]
68struct TickerState {
69 active: Vec<Arc<RebindableBlockEntityTicker>>,
70 pending: Vec<Arc<RebindableBlockEntityTicker>>,
71 by_pos: FxHashMap<BlockPos, Arc<RebindableBlockEntityTicker>>,
72 ticking: bool,
73}
74
75#[derive(Default)]
77pub(crate) struct WorldBlockEntityTickers {
78 state: SyncMutex<TickerState>,
79}
80
81impl WorldBlockEntityTickers {
82 #[must_use]
83 pub(crate) fn new() -> Self {
84 Self::default()
85 }
86
87 pub(crate) fn reconcile(
92 &self,
93 holder: &Arc<ChunkHolder>,
94 entity: SharedBlockEntity,
95 ticker: Option<BlockEntityTicker>,
96 ) {
97 let pos = entity.get_block_pos();
98 let mut state = self.state.lock();
99
100 let existing = state.by_pos.get(&pos).cloned();
101 let Some(ticker) = ticker else {
102 if let Some(wrapper) = existing
103 && wrapper
104 .binding
105 .load()
106 .as_ref()
107 .is_some_and(|binding| binding.belongs_to(holder))
108 {
109 state.by_pos.remove(&pos);
110 wrapper.rebind(None);
111 }
112 return;
113 };
114
115 let binding = Arc::new(BoundBlockEntityTicker::new(holder, entity, ticker));
116 if let Some(wrapper) = existing {
117 let same_holder = wrapper
118 .binding
119 .load()
120 .as_ref()
121 .is_some_and(|current| current.belongs_to(holder));
122 if same_holder {
123 wrapper.rebind(Some(binding));
124 return;
125 }
126 wrapper.rebind(None);
127 state.by_pos.remove(&pos);
128 }
129
130 let wrapper = Arc::new(RebindableBlockEntityTicker::new(binding));
131 state.by_pos.insert(pos, Arc::clone(&wrapper));
132 if state.ticking {
133 state.pending.push(wrapper);
134 } else {
135 state.active.push(wrapper);
136 }
137 }
138
139 pub(crate) fn remove(&self, holder: &Arc<ChunkHolder>, pos: BlockPos) {
141 let mut state = self.state.lock();
142 let Some(wrapper) = state.by_pos.get(&pos).cloned() else {
143 return;
144 };
145 if !wrapper
146 .binding
147 .load()
148 .as_ref()
149 .is_some_and(|binding| binding.belongs_to(holder))
150 {
151 return;
152 }
153 state.by_pos.remove(&pos);
154 wrapper.rebind(None);
155 }
156
157 pub(crate) fn remove_positions(&self, holder: &Arc<ChunkHolder>, positions: &[BlockPos]) {
159 let mut state = self.state.lock();
160 for pos in positions {
161 let belongs = state.by_pos.get(pos).is_some_and(|wrapper| {
162 wrapper
163 .binding
164 .load()
165 .as_ref()
166 .is_some_and(|binding| binding.belongs_to(holder))
167 });
168 if belongs && let Some(wrapper) = state.by_pos.remove(pos) {
169 wrapper.rebind(None);
170 }
171 }
172 }
173
174 pub(crate) fn tick(&self, world: &Arc<World>, runs_normally: bool) {
176 let border = runs_normally.then(|| world.world_border_snapshot());
177 self.tick_phase(runs_normally, |binding| {
178 if let Some(border) = border {
179 Self::tick_binding(world, border, binding);
180 }
181 });
182 }
183
184 fn tick_phase(
185 &self,
186 runs_normally: bool,
187 mut tick_binding: impl FnMut(&Arc<BoundBlockEntityTicker>),
188 ) {
189 let mut current = {
190 let mut state = self.state.lock();
191 debug_assert!(!state.ticking, "block-entity ticker phase re-entered");
192 state.ticking = true;
193 let pending = mem::take(&mut state.pending);
194 state.active.extend(pending);
195 mem::take(&mut state.active)
196 };
197
198 current.retain(|wrapper| {
199 let binding_guard = wrapper.binding.load();
200 let Some(binding) = binding_guard.as_ref() else {
201 return false;
202 };
203 if binding.entity.is_removed() {
204 self.remove_if_same(wrapper, binding);
205 return false;
206 }
207 if runs_normally {
208 tick_binding(binding);
209 }
210 drop(binding_guard);
211 wrapper.binding.load().is_some()
212 });
213
214 let mut state = self.state.lock();
215 debug_assert!(state.active.is_empty());
216 state.active = current;
217 state.ticking = false;
218 }
219
220 fn tick_binding(
221 world: &Arc<World>,
222 border: WorldBorderSnapshot,
223 binding: &Arc<BoundBlockEntityTicker>,
224 ) {
225 let Some(holder) = binding.holder.upgrade() else {
226 return;
227 };
228 if !holder
229 .simulation_level()
230 .is_some_and(ChunkTicketLevel::is_block_ticking)
231 || !holder.ticking_readiness_snapshot().is_block_ticking()
232 {
233 return;
234 }
235 let pos = binding.entity.get_block_pos();
236 if !border.is_block_within_bounds(pos)
237 || !world.entity_manager().is_chunk_loaded(holder.get_pos())
238 {
239 return;
240 }
241
242 let Some(state) =
243 world
244 .chunk_map
245 .block_entity_tick_state_if_owned(&holder, pos, &binding.entity)
246 else {
247 return;
248 };
249
250 if !binding.entity.is_valid_block_state(state) {
251 if !binding.logged_invalid_state.swap(true, Ordering::Relaxed) {
252 tracing::warn!(
253 block_entity_type = %binding.entity.get_type().key,
254 ?pos,
255 block = %state.get_block().key,
256 "Block entity has an invalid state for ticking"
257 );
258 }
259 return;
260 }
261
262 binding.logged_invalid_state.store(false, Ordering::Relaxed);
263 binding
264 .ticker
265 .tick(world, pos, state, binding.entity.as_ref());
266 }
267
268 fn remove_if_same(
269 &self,
270 wrapper: &Arc<RebindableBlockEntityTicker>,
271 binding: &Arc<BoundBlockEntityTicker>,
272 ) {
273 let pos = binding.entity.get_block_pos();
274 let mut state = self.state.lock();
275 let still_same_wrapper = state
276 .by_pos
277 .get(&pos)
278 .is_some_and(|current| Arc::ptr_eq(current, wrapper));
279 let still_same_binding = wrapper
280 .binding
281 .load()
282 .as_ref()
283 .is_some_and(|current| Arc::ptr_eq(current, binding));
284 if still_same_wrapper && still_same_binding {
285 state.by_pos.remove(&pos);
286 wrapper.rebind(None);
287 }
288 }
289
290 #[cfg(test)]
291 pub(crate) fn registered_len(&self) -> usize {
292 self.state.lock().by_pos.len()
293 }
294
295 #[cfg(test)]
296 pub(crate) fn active_positions(&self) -> Vec<BlockPos> {
297 self.state
298 .lock()
299 .active
300 .iter()
301 .filter_map(|wrapper| {
302 wrapper
303 .binding
304 .load()
305 .as_ref()
306 .map(|binding| binding.entity.get_block_pos())
307 })
308 .collect()
309 }
310}
311
312#[cfg(test)]
313mod tests {
314 use std::sync::{Arc, Weak};
315
316 use steel_registry::{init_vanilla_registry, vanilla_block_entity_types, vanilla_blocks};
317 use steel_utils::{BlockPos, ChunkPos};
318
319 use super::*;
320 use crate::{
321 block_entity::entities::SignBlockEntity,
322 chunk::{chunk_holder::ChunkHolder, chunk_ticket_manager::ChunkTicketLevel},
323 };
324
325 fn holder() -> Arc<ChunkHolder> {
326 Arc::new(ChunkHolder::new(
327 ChunkPos::new(0, 0),
328 ChunkTicketLevel::BLOCK_TICKING_CHUNK,
329 Some(ChunkTicketLevel::BLOCK_TICKING_CHUNK),
330 -64,
331 384,
332 ))
333 }
334
335 fn sign(pos: BlockPos) -> SharedBlockEntity {
336 Arc::new(SignBlockEntity::new(
337 Weak::new(),
338 pos,
339 vanilla_blocks::OAK_SIGN.default_state(),
340 ))
341 }
342
343 fn sign_ticker() -> BlockEntityTicker {
344 BlockEntityTicker::for_entity_tick(&vanilla_block_entity_types::SIGN)
345 }
346
347 #[test]
348 fn additions_during_phase_wait_and_follow_between_phase_additions() {
349 init_vanilla_registry();
350 let manager = WorldBlockEntityTickers::new();
351 let holder = holder();
352 let first = sign(BlockPos::new(1, 2, 3));
353 let pending = sign(BlockPos::new(2, 2, 3));
354 let between = sign(BlockPos::new(3, 2, 3));
355 manager.reconcile(&holder, Arc::clone(&first), Some(sign_ticker()));
356
357 let mut added = false;
358 manager.tick_phase(true, |_| {
359 if !added {
360 added = true;
361 manager.reconcile(&holder, Arc::clone(&pending), Some(sign_ticker()));
362 }
363 });
364 manager.reconcile(&holder, Arc::clone(&between), Some(sign_ticker()));
365
366 let mut observed = Vec::new();
367 manager.tick_phase(true, |binding| {
368 observed.push(binding.entity.get_block_pos());
369 });
370 assert_eq!(
371 observed,
372 [
373 first.get_block_pos(),
374 between.get_block_pos(),
375 pending.get_block_pos(),
376 ]
377 );
378 }
379
380 #[test]
381 fn rebind_before_turn_uses_new_owner_in_the_original_slot() {
382 init_vanilla_registry();
383 let manager = WorldBlockEntityTickers::new();
384 let holder = holder();
385 let first = sign(BlockPos::new(1, 2, 3));
386 let old_second = sign(BlockPos::new(2, 2, 3));
387 let new_second = sign(BlockPos::new(2, 2, 3));
388 manager.reconcile(&holder, Arc::clone(&first), Some(sign_ticker()));
389 manager.reconcile(&holder, Arc::clone(&old_second), Some(sign_ticker()));
390
391 let mut observed = Vec::new();
392 manager.tick_phase(true, |binding| {
393 if Arc::ptr_eq(&binding.entity, &first) {
394 manager.reconcile(&holder, Arc::clone(&new_second), Some(sign_ticker()));
395 observed.push("first");
396 } else if Arc::ptr_eq(&binding.entity, &new_second) {
397 observed.push("new");
398 } else {
399 observed.push("old");
400 }
401 });
402
403 assert_eq!(observed, ["first", "new"]);
404 assert_eq!(manager.registered_len(), 2);
405 }
406
407 #[test]
408 fn remove_then_add_during_phase_creates_a_pending_tail() {
409 init_vanilla_registry();
410 let manager = WorldBlockEntityTickers::new();
411 let holder = holder();
412 let first = sign(BlockPos::new(1, 2, 3));
413 let old_second = sign(BlockPos::new(2, 2, 3));
414 let new_second = sign(BlockPos::new(2, 2, 3));
415 manager.reconcile(&holder, Arc::clone(&first), Some(sign_ticker()));
416 manager.reconcile(&holder, Arc::clone(&old_second), Some(sign_ticker()));
417
418 let mut changed = false;
419 let mut first_phase = Vec::new();
420 manager.tick_phase(true, |binding| {
421 first_phase.push(Arc::clone(&binding.entity));
422 if !changed && Arc::ptr_eq(&binding.entity, &first) {
423 changed = true;
424 manager.remove(&holder, old_second.get_block_pos());
425 manager.reconcile(&holder, Arc::clone(&new_second), Some(sign_ticker()));
426 }
427 });
428 assert_eq!(first_phase.len(), 1);
429 assert!(Arc::ptr_eq(&first_phase[0], &first));
430
431 let mut second_phase = Vec::new();
432 manager.tick_phase(true, |binding| {
433 second_phase.push(Arc::clone(&binding.entity));
434 });
435 assert_eq!(second_phase.len(), 2);
436 assert!(Arc::ptr_eq(&second_phase[0], &first));
437 assert!(Arc::ptr_eq(&second_phase[1], &new_second));
438 }
439
440 #[test]
441 fn frozen_phase_merges_pending_and_prunes_unbound_without_callbacks() {
442 init_vanilla_registry();
443 let manager = WorldBlockEntityTickers::new();
444 let holder = holder();
445 let first = sign(BlockPos::new(1, 2, 3));
446 let pending = sign(BlockPos::new(2, 2, 3));
447 manager.reconcile(&holder, Arc::clone(&first), Some(sign_ticker()));
448
449 manager.tick_phase(true, |_| {
450 manager.reconcile(&holder, Arc::clone(&pending), Some(sign_ticker()));
451 manager.remove(&holder, first.get_block_pos());
452 });
453 manager.tick_phase(false, |_| panic!("frozen phase must suppress callbacks"));
454
455 let mut observed = Vec::new();
456 manager.tick_phase(true, |binding| {
457 observed.push(Arc::clone(&binding.entity));
458 });
459 assert_eq!(observed.len(), 1);
460 assert!(Arc::ptr_eq(&observed[0], &pending));
461 }
462}