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(async move {
161                        map_clone.save_chunk(&holder_clone, save_dependency).await;
162                    });
163                    return true;
164                }
165
166                let has_chunk = holder.try_chunk(ChunkStatus::Empty).is_some();
167                finalized.push((*pos, Arc::clone(holder), has_chunk));
168                false
169            });
170        }
171
172        let world = self.world_gen_context.world();
173        for (pos, holder, has_chunk) in finalized {
174            let cleared = if has_chunk {
175                holder.try_chunk(ChunkStatus::Empty).map_or_else(
176                    ClearedBlockEntities::default,
177                    |chunk| {
178                        if let Some(full) = holder.try_full_chunk() {
179                            full.clear_all_block_entities_staged()
180                        } else {
181                            chunk.clear_all_block_entities();
182                            ClearedBlockEntities::default()
183                        }
184                    },
185                )
186            } else {
187                ClearedBlockEntities::default()
188            };
189            self.finalized_block_entity_unloads
190                .lock()
191                .push(FinalizedBlockEntityUnload {
192                    holder: Arc::clone(&holder),
193                    lifecycle_dispatchers: cleared.lifecycle_dispatchers,
194                    positions: cleared.positions,
195                });
196
197            world.unregister_full_chunk_ticks(pos);
198            world.on_entity_chunk_unload_finalized(pos);
199            if has_chunk {
200                let map_clone = Arc::clone(self);
201                self.task_tracker.spawn(async move {
202                    if let Err(e) = map_clone.storage.release_chunk(pos).await {
203                        tracing::error!(?pos, "Error releasing chunk: {e}");
204                    }
205                });
206            }
207        }
208    }
209
210    /// Saves all dirty chunks to disk.
211    ///
212    /// This method should be called during graceful shutdown to ensure all
213    /// modified chunks are persisted. It saves:
214    /// 1. All dirty chunks in the active `chunks` map
215    /// 2. All chunks pending unload in the `unloading_chunks` map
216    /// 3. Closes all region file handles (flushing headers)
217    ///
218    /// Returns the number of chunks saved.
219    #[instrument(level = "info", skip(self), name = "save_all_chunks")]
220    #[expect(
221        clippy::too_many_lines,
222        reason = "shutdown persistence keeps chunk coverage and entity auditing in one pass"
223    )]
224    pub async fn save_all_chunks(self: &Arc<Self>) -> io::Result<usize> {
225        let mut saved_count = 0;
226
227        self.flush_queued_light_changes_for_save().await;
228
229        // Collect all chunks from both maps
230        let all_chunks: Vec<Arc<ChunkHolder>> = {
231            let mut chunks = Vec::new();
232            self.chunks.iter_sync(|_, holder| {
233                chunks.push(holder.clone());
234                true
235            });
236            self.unloading_chunks.iter_sync(|_, holder| {
237                chunks.push(holder.clone());
238                true
239            });
240            chunks
241        };
242        let mut covered_chunk_positions = FxHashSet::default();
243
244        tracing::info!(chunk_count = all_chunks.len(), "Saving chunks");
245
246        // Save all chunks that have data
247        for holder in &all_chunks {
248            let chunk_pos = holder.get_pos();
249            let prepared = {
250                let Some(chunk) = holder.try_chunk(ChunkStatus::StructureStarts) else {
251                    // Matches save_chunk: StructureStarts is the first persisted
252                    // chunk status, so lower-status chunks do not own durable
253                    // runtime entity data.
254                    continue;
255                };
256                let Some(status) = holder.published_status() else {
257                    continue;
258                };
259                let world = self.world_gen_context.world();
260                let runtime_entities = world
261                    .entity_manager()
262                    .get_saveable_entities_for_chunk(chunk_pos);
263                let force = world.entity_manager().has_save_pending_for_chunk(chunk_pos);
264                let dirty = chunk.take_dirty();
265                let prepared = if dirty || force {
266                    ChunkStorage::prepare_chunk_save(chunk, status, &runtime_entities, true)
267                } else {
268                    None
269                };
270                let Some(prepared) = prepared else {
271                    if dirty {
272                        chunk.mark_dirty();
273                    } else if !force {
274                        covered_chunk_positions.insert(chunk_pos);
275                    }
276                    continue;
277                };
278                prepared
279            };
280
281            let mut prepared = prepared;
282            let handled_runtime_entity_ids = mem::take(&mut prepared.handled_runtime_entity_ids);
283            let world = self.world_gen_context.world();
284            let _save_dependency = holder.add_save_dependency();
285            match self
286                .storage
287                .save_chunk_data(prepared, self.chunk_encoding_pool.as_ref())
288                .await
289            {
290                Ok(true) => {
291                    world
292                        .entity_manager()
293                        .on_chunk_saved(chunk_pos, &handled_runtime_entity_ids);
294                    covered_chunk_positions.insert(chunk_pos);
295                    saved_count += 1;
296                }
297                Ok(false) => Self::mark_chunk_dirty_for_save_retry(holder),
298                Err(e) => {
299                    tracing::error!(chunk = ?holder.get_pos(), "Failed to save chunk: {e}");
300                    Self::mark_chunk_dirty_for_save_retry(holder);
301                }
302            }
303        }
304
305        let world = self.world_gen_context.world();
306        let covered_chunk_positions = covered_chunk_positions.into_iter().collect::<Vec<_>>();
307        let unsaved_entities = world
308            .entity_manager()
309            .saveable_entities_outside_chunks(&covered_chunk_positions);
310        if !unsaved_entities.is_empty() {
311            let chunk_count = unsaved_entities
312                .iter()
313                .map(|entity| entity.chunk)
314                .collect::<FxHashSet<_>>()
315                .len();
316            let sample = unsaved_entities
317                .iter()
318                .take(16)
319                .map(|entity| format!("{}:{}@{:?}", entity.entity_id, entity.uuid, entity.chunk))
320                .collect::<Vec<_>>()
321                .join(", ");
322            tracing::warn!(
323                entity_count = unsaved_entities.len(),
324                chunk_count,
325                sample = %sample,
326                "Saveable runtime entities remain in chunks without save holders after chunk save"
327            );
328        }
329
330        // Close all region files (flushes headers and releases file handles)
331        if let Err(e) = self.storage.close_all().await {
332            tracing::error!("Failed to close region files: {e}");
333        }
334
335        tracing::info!(
336            saved_count,
337            total_checked = all_chunks.len(),
338            "Chunk save complete"
339        );
340
341        Ok(saved_count)
342    }
343}