Skip to main content

steel_core/chunk/chunk_map/
persistence.rs

1use super::{
2    Arc, Chunk, ChunkHolder, ChunkMap, ChunkPos, ChunkSaveDependency, ChunkStatus, ChunkStorage,
3    ClearedBlockEntities, FinalizedBlockEntityUnload, FxHashSet, instrument, io, mem,
4};
5use crate::chunk_saver::PreparedChunkSave;
6use tokio::sync::oneshot;
7
8impl ChunkMap {
9    async fn prepare_chunk_save_on_pool(
10        self: &Arc<Self>,
11        chunk_holder: &Arc<ChunkHolder>,
12    ) -> Option<PreparedChunkSave> {
13        let (sender, receiver) = oneshot::channel();
14        let map = Arc::clone(self);
15        let holder = Arc::clone(chunk_holder);
16        self.chunk_encoding_pool.spawn(move || {
17            let chunk_pos = holder.get_pos();
18            let prepared = {
19                let Some(save_preparation) = holder.try_begin_save_preparation() else {
20                    let _ = sender.send(None);
21                    return;
22                };
23                let Some(chunk_guard) = holder.try_chunk(ChunkStatus::StructureStarts) else {
24                    // Vanilla only persists chunks once they reach StructureStarts.
25                    // Runtime entities in lower-status chunks are an accepted loss
26                    // on unload/shutdown until those chunks cross that boundary.
27                    let _ = sender.send(None);
28                    return;
29                };
30
31                let Some(status) = holder.published_status() else {
32                    let _ = sender.send(None);
33                    return;
34                };
35
36                let world = map.world_gen_context.world();
37                let runtime_entities = world
38                    .entity_manager()
39                    .get_saveable_entities_for_chunk(chunk_pos);
40                let force = world.entity_manager().has_save_pending_for_chunk(chunk_pos);
41                let dirty = chunk_guard.take_dirty();
42                let prepared = if dirty || force {
43                    ChunkStorage::prepare_chunk_save(chunk_guard, status, &runtime_entities, true)
44                } else {
45                    None
46                };
47
48                if prepared.is_none() && dirty {
49                    chunk_guard.mark_dirty();
50                }
51
52                // Revival need not wait for encoding or disk I/O once this owned input exists.
53                drop(save_preparation);
54                prepared
55            };
56
57            if let Err(unsent) = sender.send(prepared) {
58                if unsent.is_some() {
59                    Self::mark_chunk_dirty_for_save_retry(&holder);
60                }
61                tracing::trace!(
62                    chunk = ?chunk_pos,
63                    "Discarding prepared chunk after its save task was canceled"
64                );
65            }
66        });
67
68        if let Ok(prepared) = receiver.await {
69            prepared
70        } else {
71            Self::mark_chunk_dirty_for_save_retry(chunk_holder);
72            tracing::error!(
73                chunk = ?chunk_holder.get_pos(),
74                "Chunk save preparation task ended without returning a result"
75            );
76            None
77        }
78    }
79
80    /// Saves a chunk to disk. Does not remove from `unloading_chunks`.
81    #[instrument(level = "trace", skip(self, chunk_holder, _save_dependency), fields(chunk = ?chunk_holder.get_pos()))]
82    pub(super) async fn save_chunk(
83        self: &Arc<Self>,
84        chunk_holder: &Arc<ChunkHolder>,
85        _save_dependency: ChunkSaveDependency,
86    ) {
87        let chunk_pos = chunk_holder.get_pos();
88        self.flush_queued_light_changes_touching_chunk_for_save(chunk_pos)
89            .await;
90
91        // Save chunk data if dirty. Preparation runs on the shared chunk encoding
92        // pool so snapshot locking and persistence conversion never block Tokio workers.
93        if let Some(mut prepared) = self.prepare_chunk_save_on_pool(chunk_holder).await {
94            let handled_runtime_entity_ids = mem::take(&mut prepared.handled_runtime_entity_ids);
95            let world = self.world_gen_context.world();
96            match self
97                .storage
98                .save_chunk_data(prepared, self.chunk_encoding_pool.as_ref())
99                .await
100            {
101                Ok(true) => world
102                    .entity_manager()
103                    .on_chunk_saved(chunk_pos, &handled_runtime_entity_ids),
104                Ok(false) => Self::mark_chunk_dirty_for_save_retry(chunk_holder),
105                Err(e) => {
106                    tracing::error!("Error saving chunk: {e}");
107                    Self::mark_chunk_dirty_for_save_retry(chunk_holder);
108                }
109            }
110        }
111    }
112
113    pub(super) fn mark_chunk_dirty_for_save_retry(chunk_holder: &ChunkHolder) {
114        let Some(chunk) = chunk_holder.try_chunk(ChunkStatus::StructureStarts) else {
115            return;
116        };
117        chunk.mark_dirty();
118    }
119
120    /// Processes chunks that are pending unload.
121    ///
122    /// Iterates over `unloading_chunks`. For each chunk with `strong_count == 1`:
123    /// - If staged to revive at the next lifecycle boundary: keep
124    /// - If dirty: spawn save task (keep until saved and clean)
125    /// - If not dirty: release region handle and remove
126    #[instrument(level = "trace", skip(self, staged_revivals))]
127    pub(super) fn process_unloads(self: &Arc<Self>, staged_revivals: &FxHashSet<ChunkPos>) {
128        self.propagate_queued_light_changes();
129
130        let mut finalized = Vec::new();
131        {
132            let light_updates = self.light_updates.lock();
133            self.unloading_chunks.retain_sync(|pos, holder| {
134                // Prepared ticket changes publish only at the next lifecycle boundary.
135                if staged_revivals.contains(pos) {
136                    return true;
137                }
138
139                if light_updates.touches_chunk(*pos) {
140                    return true;
141                }
142
143                if Arc::strong_count(holder) != 1 {
144                    return true;
145                }
146
147                let is_dirty = holder
148                    .try_chunk(ChunkStatus::StructureStarts)
149                    .is_some_and(Chunk::is_dirty);
150                let has_save_pending_entities = self
151                    .world_gen_context
152                    .world()
153                    .entity_manager()
154                    .has_save_pending_for_chunk(*pos);
155
156                if is_dirty || has_save_pending_entities {
157                    let save_dependency = holder.add_save_dependency();
158                    let holder_clone = Arc::clone(holder);
159                    let map_clone = Arc::clone(self);
160                    self.task_tracker.spawn_on(
161                        async move {
162                            map_clone.save_chunk(&holder_clone, save_dependency).await;
163                        },
164                        self.chunk_runtime.handle(),
165                    );
166                    return true;
167                }
168
169                let has_chunk = holder.try_chunk(ChunkStatus::Empty).is_some();
170                finalized.push((*pos, Arc::clone(holder), has_chunk));
171                false
172            });
173        }
174
175        let world = self.world_gen_context.world();
176        for (pos, holder, has_chunk) in finalized {
177            let cleared = if has_chunk {
178                holder.try_chunk(ChunkStatus::Empty).map_or_else(
179                    ClearedBlockEntities::default,
180                    |chunk| {
181                        if let Some(full) = holder.try_full_chunk() {
182                            full.clear_all_block_entities_staged()
183                        } else {
184                            chunk.clear_all_block_entities();
185                            ClearedBlockEntities::default()
186                        }
187                    },
188                )
189            } else {
190                ClearedBlockEntities::default()
191            };
192            self.finalized_block_entity_unloads
193                .lock()
194                .push(FinalizedBlockEntityUnload {
195                    holder: Arc::clone(&holder),
196                    lifecycle_dispatchers: cleared.lifecycle_dispatchers,
197                    positions: cleared.positions,
198                });
199
200            world.unregister_full_chunk_ticks(pos);
201            world.on_entity_chunk_unload_finalized(pos);
202            if has_chunk {
203                let map_clone = Arc::clone(self);
204                self.task_tracker.spawn_on(
205                    async move {
206                        if let Err(e) = map_clone.storage.release_chunk(pos).await {
207                            tracing::error!(?pos, "Error releasing chunk: {e}");
208                        }
209                    },
210                    self.chunk_runtime.handle(),
211                );
212            }
213        }
214    }
215
216    /// Saves all dirty chunks to disk.
217    ///
218    /// This method should be called during graceful shutdown to ensure all
219    /// modified chunks are persisted. It saves:
220    /// 1. All dirty chunks in the active `chunks` map
221    /// 2. All chunks pending unload in the `unloading_chunks` map
222    /// 3. Closes all region file handles (flushing headers)
223    ///
224    /// Returns the number of chunks saved.
225    #[instrument(level = "info", skip(self), name = "save_all_chunks")]
226    #[expect(
227        clippy::too_many_lines,
228        reason = "shutdown persistence keeps chunk coverage and entity auditing in one pass"
229    )]
230    pub async fn save_all_chunks(self: &Arc<Self>) -> io::Result<usize> {
231        let mut saved_count = 0;
232
233        self.flush_queued_light_changes_for_save().await;
234
235        // Collect all chunks from both maps
236        let all_chunks: Vec<Arc<ChunkHolder>> = {
237            let mut chunks = Vec::new();
238            self.chunks.iter_sync(|_, holder| {
239                chunks.push(holder.clone());
240                true
241            });
242            self.unloading_chunks.iter_sync(|_, holder| {
243                chunks.push(holder.clone());
244                true
245            });
246            chunks
247        };
248        let mut covered_chunk_positions = FxHashSet::default();
249
250        tracing::info!(chunk_count = all_chunks.len(), "Saving chunks");
251
252        // Save all chunks that have data
253        for holder in &all_chunks {
254            let chunk_pos = holder.get_pos();
255            let prepared = {
256                let Some(chunk) = holder.try_chunk(ChunkStatus::StructureStarts) else {
257                    // Matches save_chunk: StructureStarts is the first persisted
258                    // chunk status, so lower-status chunks do not own durable
259                    // runtime entity data.
260                    continue;
261                };
262                let Some(status) = holder.published_status() else {
263                    continue;
264                };
265                let world = self.world_gen_context.world();
266                let runtime_entities = world
267                    .entity_manager()
268                    .get_saveable_entities_for_chunk(chunk_pos);
269                let force = world.entity_manager().has_save_pending_for_chunk(chunk_pos);
270                let dirty = chunk.take_dirty();
271                let prepared = if dirty || force {
272                    ChunkStorage::prepare_chunk_save(chunk, status, &runtime_entities, true)
273                } else {
274                    None
275                };
276                let Some(prepared) = prepared else {
277                    if dirty {
278                        chunk.mark_dirty();
279                    } else if !force {
280                        covered_chunk_positions.insert(chunk_pos);
281                    }
282                    continue;
283                };
284                prepared
285            };
286
287            let mut prepared = prepared;
288            let handled_runtime_entity_ids = mem::take(&mut prepared.handled_runtime_entity_ids);
289            let world = self.world_gen_context.world();
290            let _save_dependency = holder.add_save_dependency();
291            match self
292                .storage
293                .save_chunk_data(prepared, self.chunk_encoding_pool.as_ref())
294                .await
295            {
296                Ok(true) => {
297                    world
298                        .entity_manager()
299                        .on_chunk_saved(chunk_pos, &handled_runtime_entity_ids);
300                    covered_chunk_positions.insert(chunk_pos);
301                    saved_count += 1;
302                }
303                Ok(false) => Self::mark_chunk_dirty_for_save_retry(holder),
304                Err(e) => {
305                    tracing::error!(chunk = ?holder.get_pos(), "Failed to save chunk: {e}");
306                    Self::mark_chunk_dirty_for_save_retry(holder);
307                }
308            }
309        }
310
311        let world = self.world_gen_context.world();
312        let covered_chunk_positions = covered_chunk_positions.into_iter().collect::<Vec<_>>();
313        let unsaved_entities = world
314            .entity_manager()
315            .saveable_entities_outside_chunks(&covered_chunk_positions);
316        if !unsaved_entities.is_empty() {
317            let chunk_count = unsaved_entities
318                .iter()
319                .map(|entity| entity.chunk)
320                .collect::<FxHashSet<_>>()
321                .len();
322            let sample = unsaved_entities
323                .iter()
324                .take(16)
325                .map(|entity| format!("{}:{}@{:?}", entity.entity_id, entity.uuid, entity.chunk))
326                .collect::<Vec<_>>()
327                .join(", ");
328            tracing::warn!(
329                entity_count = unsaved_entities.len(),
330                chunk_count,
331                sample = %sample,
332                "Saveable runtime entities remain in chunks without save holders after chunk save"
333            );
334        }
335
336        // Close all region files (flushes headers and releases file handles)
337        if let Err(e) = self.storage.close_all().await {
338            tracing::error!("Failed to close region files: {e}");
339        }
340
341        tracing::info!(
342            saved_count,
343            total_checked = all_chunks.len(),
344            "Chunk save complete"
345        );
346
347        Ok(saved_count)
348    }
349}