Skip to main content

fxfs/object_store/
flush.rs

1// Copyright 2021 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5// This module is responsible for flushing (a.k.a. compacting) the object store trees.
6
7use crate::errors::FxfsError;
8use crate::filesystem::{FlushReason, ForceMajor};
9use crate::log::*;
10use crate::lsm_tree::types::{ItemRef, LayerIterator};
11use crate::lsm_tree::{LSMTree, layer_from_handle};
12use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, ReadObjectHandle};
13use crate::object_store::extent_record::ExtentValue;
14use crate::object_store::object_manager::{ObjectManager, ReservationUpdate};
15use crate::object_store::object_record::{ObjectKey, ObjectValue};
16use crate::object_store::transaction::{AssociatedObject, LockKey, Mutation, lock_keys};
17use crate::object_store::{
18    AssocObj, DirectWriter, EncryptedMutations, HandleOptions, LastObjectId, LastObjectIdInfo,
19    LockState, MAX_ENCRYPTED_MUTATIONS_SIZE, ObjectStore, Options, ReservedId, StoreInfo,
20    layer_size_from_encrypted_mutations_size, tree,
21};
22use crate::serialized_types::{LATEST_VERSION, Version, VersionedLatest};
23use anyhow::{Context, Error, anyhow};
24use fxfs_crypto::UnwrappedKey;
25use std::sync::OnceLock;
26use std::sync::atomic::Ordering;
27
28#[fxfs_trace::trace]
29impl ObjectStore {
30    /// Takes a flush lock on self and performs the flush with a default
31    /// FlushReason::Journal(ForceMajor::False).
32    pub async fn flush(&self) -> Result<Version, Error> {
33        self.flush_with_reason(FlushReason::Journal(ForceMajor::False)).await
34    }
35
36    /// Takes a flush lock on self and performs the flush.
37    pub async fn flush_with_reason(&self, reason: FlushReason) -> Result<Version, Error> {
38        let filesystem = self.filesystem();
39
40        let keys = lock_keys![LockKey::flush(self.store_object_id())];
41        let _guard = filesystem.lock_manager().write_lock(keys).await;
42
43        self.flush_guarded_with_reason(reason).await
44    }
45
46    /// Performs a flush while already holding a flush guard on self.
47    #[trace("store_object_id" => self.store_object_id)]
48    pub async fn flush_guarded_with_reason(&self, reason: FlushReason) -> Result<Version, Error> {
49        if self.parent_store.is_none() {
50            // Early exit, but still return the earliest version used by a struct in the tree
51            return Ok(self.tree.get_earliest_version());
52        }
53
54        // After taking the lock, check to see if the store has been deleted.
55        if matches!(*self.lock_state.lock(), LockState::Deleted) {
56            // When we compact, it's possible that the store has been deleted since we gathered the
57            // list of stores that need compacting.  This is benign.
58            return Ok(LATEST_VERSION);
59        }
60
61        let filesystem = self.filesystem();
62        let object_manager = filesystem.object_manager();
63        let earliest_version = self.tree.get_earliest_version();
64        let needs_flush = object_manager.needs_flush(self.store_object_id);
65        // If we don't need to do anything for the stated flush purpose, do nothing.
66        if (reason == FlushReason::Journal(ForceMajor::False) && !needs_flush)
67            || (reason == FlushReason::UpgradeVersion && earliest_version == LATEST_VERSION)
68            || (reason == FlushReason::EncryptedMutations
69                && self.store_info().unwrap().encrypted_mutations_object_id == INVALID_OBJECT_ID)
70        {
71            // Early exit, but still return the earliest version used by a struct in the
72            // tree.
73            return Ok(earliest_version);
74        }
75
76        let trace = self.trace.load(Ordering::Relaxed);
77        if trace {
78            info!(store_id = self.store_object_id(); "OS: begin flush");
79        }
80
81        if matches!(&*self.lock_state.lock(), LockState::Locked) {
82            self.flush_locked().await.with_context(|| {
83                format!("Failed to flush object store {}", self.store_object_id)
84            })?;
85        } else {
86            self.flush_unlocked(reason).await.with_context(|| {
87                format!("Failed to flush object store {}", self.store_object_id)
88            })?;
89        }
90
91        if trace {
92            info!(store_id = self.store_object_id(); "OS: end flush");
93        }
94        if let Some(callback) = &*self.flush_callback.lock() {
95            callback(self);
96        }
97
98        // NOTE: `num_flushes` must only be incremented while the flush lock is held, as
99        // `ObjectStore::unlock` relies on this to detect concurrent flushes.
100        let mut counters = self.counters.lock();
101        counters.num_flushes += 1;
102        counters.last_flush_time = Some(std::time::SystemTime::now());
103        // Return the earliest version used by a struct in the tree
104        Ok(self.tree.get_earliest_version())
105    }
106
107    // Flushes an unlocked store. Returns the layer file sizes.
108    async fn flush_unlocked(&self, reason: FlushReason) -> Result<Vec<u64>, Error> {
109        struct StoreInfoSnapshot<'a> {
110            store: &'a ObjectStore,
111            store_info: OnceLock<StoreInfo>,
112        }
113        impl AssociatedObject for StoreInfoSnapshot<'_> {
114            fn will_apply_mutation(
115                &self,
116                _mutation: &Mutation,
117                _object_id: u64,
118                _manager: &ObjectManager,
119            ) {
120                let mut store_info = self.store.store_info().unwrap();
121
122                let lock_state = self.store.lock_state.lock();
123                if let LockState::Unlocked { mutations_cipher, .. } = &*lock_state {
124                    store_info.mutations_cipher_offset = mutations_cipher.sequence_number();
125                }
126
127                self.store_info.set(store_info).unwrap();
128            }
129        }
130
131        let store_info_snapshot = StoreInfoSnapshot { store: self, store_info: OnceLock::new() };
132
133        let filesystem = self.filesystem();
134        let object_manager = filesystem.object_manager();
135        let reservation = object_manager.metadata_reservation();
136        let txn_options = Options {
137            skip_journal_checks: true,
138            borrow_metadata_space: true,
139            allocator_reservation: Some(reservation),
140            ..Default::default()
141        };
142
143        // The BeginFlush mutation must be within a transaction that has no impact on StoreInfo
144        // since we want to get an accurate snapshot of StoreInfo.
145        let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
146        transaction.add_with_object(
147            self.store_object_id(),
148            Mutation::BeginFlush,
149            AssocObj::Borrowed(&store_info_snapshot),
150        );
151        transaction.commit().await.context("Failed to commit BeginFlush transaction")?;
152
153        let mut new_store_info = store_info_snapshot.store_info.into_inner().unwrap();
154
155        // There is a transaction to create objects at the start and then another transaction at the
156        // end. Between those two transactions, there are transactions that write to the files.  In
157        // the first transaction, objects are created in the graveyard. Upon success, the objects
158        // are removed from the graveyard.
159        let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
160
161        // Create and write a new layer, compacting existing layers.
162        let parent_store = self.parent_store.as_ref().unwrap();
163        let handle_options = HandleOptions { skip_journal_checks: true, ..Default::default() };
164        let id_and_key = {
165            let mut lock_state = self.lock_state.lock();
166            match &mut *lock_state {
167                LockState::Unlocked { cached_keys, .. } => {
168                    if let Some(item) = cached_keys.pop() {
169                        Ok(Some(item))
170                    } else {
171                        log::warn!("No cached keys for flush for store {}", self.store_object_id());
172                        Err(anyhow!(FxfsError::Internal).context("No cached keys for flush"))
173                    }
174                }
175                LockState::UnlockedReadOnly(..) => {
176                    Err(anyhow!(FxfsError::Internal).context("Flush on read-only store"))
177                }
178                LockState::Unencrypted => Ok(None),
179                _ => Err(anyhow!(FxfsError::Internal))
180                    .with_context(|| format!("Invalid lock state ({:?}) for flush", *lock_state)),
181            }
182        }?;
183
184        let (new_object_tree_layer, unwrapped_key) = if let Some((raw_id, key, unwrapped_key)) =
185            id_and_key
186        {
187            let object_id = ReservedId::new(parent_store, raw_id);
188            let handle = ObjectStore::create_object_with_key(
189                parent_store,
190                &mut transaction,
191                object_id,
192                handle_options,
193                key,
194                UnwrappedKey::new(unwrapped_key.clone()),
195            )
196            .await?;
197            (handle, Some(unwrapped_key))
198        } else {
199            let handle =
200                ObjectStore::create_object(parent_store, &mut transaction, handle_options, None)
201                    .await?;
202            (handle, None)
203        };
204        let writer = DirectWriter::new(&new_object_tree_layer, txn_options).await;
205        let new_object_tree_layer_object_id = new_object_tree_layer.object_id();
206        parent_store.add_to_graveyard(&mut transaction, new_object_tree_layer_object_id);
207
208        transaction.commit().await.context("Failed to commit create layer transaction")?;
209
210        // *Do* the actual compaction.
211        let full_compaction = match reason {
212            FlushReason::UpgradeVersion | FlushReason::Journal(ForceMajor::True) => true,
213            _ => false,
214        };
215        let (layers_to_keep, old_layers) = tree::flush(
216            &self.tree,
217            writer,
218            matches!(reason, FlushReason::Journal(_))
219                .then(|| filesystem.journal().get_compaction_yielder()),
220            full_compaction,
221        )
222        .await
223        .context("Failed to flush tree")?;
224
225        // Finalise the compaction.
226        let mut new_layers = vec![layer_from_handle(new_object_tree_layer, unwrapped_key).await?];
227        new_layers.extend(layers_to_keep.iter().map(|l| (*l).clone()));
228
229        new_store_info.layers = Vec::new();
230        for layer in &new_layers {
231            if let Some(handle) = layer.handle() {
232                new_store_info.layers.push(handle.object_id());
233            }
234        }
235
236        let reservation_update: ReservationUpdate; // Must live longer than end_transaction.
237        let mut end_transaction = parent_store
238            .new_transaction(
239                lock_keys![LockKey::object(
240                    self.parent_store.as_ref().unwrap().store_object_id(),
241                    self.store_info_handle_object_id().unwrap(),
242                )],
243                txn_options,
244            )
245            .await?;
246
247        parent_store.remove_from_graveyard(&mut end_transaction, new_object_tree_layer_object_id);
248
249        // Move the existing layers we're compacting to the graveyard at the end.
250        for layer in &old_layers {
251            if let Some(handle) = layer.handle() {
252                parent_store.add_to_graveyard(&mut end_transaction, handle.object_id());
253            }
254        }
255
256        let old_encrypted_mutations_object_id =
257            std::mem::replace(&mut new_store_info.encrypted_mutations_object_id, INVALID_OBJECT_ID);
258        if old_encrypted_mutations_object_id != INVALID_OBJECT_ID {
259            parent_store.add_to_graveyard(&mut end_transaction, old_encrypted_mutations_object_id);
260        }
261
262        // `last_object_id` is updated differently to other members of `StoreInfo`.  We must ensure
263        // that those fields match the current in-memory values.  See the lengthy comment in
264        // `get_next_object_id` for more information.  `end_transaction` has a lock on the same lock
265        // that `get_next_object_id` uses, so there's no danger of the key changing now.
266        //
267        // This might capture object IDs that might be in transactions not yet committed.  In
268        // theory, we could do better than this but it's not worth the effort.
269        match &mut new_store_info.last_object_id {
270            LastObjectIdInfo::Unencrypted { id } => {
271                let LastObjectId::Unencrypted { id: in_memory_value } =
272                    &*self.last_object_id.lock()
273                else {
274                    unreachable!()
275                };
276                *id = *in_memory_value;
277            }
278            LastObjectIdInfo::Encrypted { id, key } => {
279                let LastObjectId::Encrypted { id: in_memory_value, .. } =
280                    &*self.last_object_id.lock()
281                else {
282                    unreachable!()
283                };
284                *id = *in_memory_value;
285                let guard = self.store_info.lock();
286                let current_store_info = guard.as_ref().unwrap();
287                let LastObjectIdInfo::Encrypted { key: in_memory_value, .. } =
288                    &current_store_info.last_object_id
289                else {
290                    unreachable!()
291                };
292                *key = in_memory_value.clone();
293            }
294            LastObjectIdInfo::Low32Bit => {}
295        }
296
297        self.write_store_info(&mut end_transaction, &new_store_info).await?;
298
299        let layer_file_sizes = new_layers
300            .iter()
301            .map(|l| l.handle().map(ReadObjectHandle::get_size).unwrap_or(0))
302            .collect::<Vec<u64>>();
303
304        let total_layer_size = layer_file_sizes.iter().sum();
305        reservation_update =
306            ReservationUpdate::new(tree::reservation_amount_from_layer_size(total_layer_size));
307
308        end_transaction.add_with_object(
309            self.store_object_id(),
310            Mutation::EndFlush,
311            AssocObj::Borrowed(&reservation_update),
312        );
313
314        if self.trace.load(Ordering::Relaxed) {
315            info!(
316                store_id = self.store_object_id(),
317                old_layer_count = old_layers.len(),
318                new_layer_count = new_layers.len(),
319                total_layer_size,
320                new_store_info:?;
321                "OS: compacting"
322            );
323        }
324
325        end_transaction
326            .commit_with_callback(|_| {
327                let mut store_info = self.store_info.lock();
328                let info = store_info.as_mut().unwrap();
329                info.layers = new_store_info.layers;
330                info.encrypted_mutations_object_id = new_store_info.encrypted_mutations_object_id;
331                info.mutations_cipher_offset = new_store_info.mutations_cipher_offset;
332                self.tree.set_layers(new_layers);
333            })
334            .await
335            .context("Failed to commit EndFlush transaction")?;
336
337        // Now close the layers and purge them.
338        for layer in old_layers {
339            let object_id = layer.handle().map(|h| h.object_id());
340            layer.close_layer().await;
341            if let Some(object_id) = object_id {
342                parent_store
343                    .tombstone_object(object_id, txn_options, None)
344                    .await
345                    .context("Failed to tombstone old layer")?;
346            }
347        }
348
349        if old_encrypted_mutations_object_id != INVALID_OBJECT_ID {
350            parent_store
351                .tombstone_object(old_encrypted_mutations_object_id, txn_options, None)
352                .await
353                .context("Failed to tombstone old encrypted mutations")?;
354        }
355
356        Ok(layer_file_sizes)
357    }
358
359    // Flushes a locked store.
360    async fn flush_locked(&self) -> Result<(), Error> {
361        let filesystem = self.filesystem();
362        let object_manager = filesystem.object_manager();
363        let reservation = object_manager.metadata_reservation();
364        let txn_options = Options {
365            skip_journal_checks: true,
366            borrow_metadata_space: true,
367            allocator_reservation: Some(reservation),
368            ..Default::default()
369        };
370
371        let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
372        transaction.add(self.store_object_id(), Mutation::BeginFlush);
373        transaction.commit().await.context("Failed to commit BeginFlush transaction")?;
374
375        let mut new_store_info = self.load_store_info().await?;
376
377        // There is a transaction to create objects at the start and then another transaction at the
378        // end. Between those two transactions, there are transactions that write to the files.  In
379        // the first transaction, objects are created in the graveyard. Upon success, the objects
380        // are removed from the graveyard.
381        let mut transaction = self.new_transaction(lock_keys![], txn_options).await?;
382
383        let reservation_update: ReservationUpdate; // Must live longer than end_transaction.
384        let handle; // Must live longer than end_transaction.
385        let mut end_transaction;
386
387        // We need to either write our encrypted mutations to a new file, or append them to an
388        // existing one.
389        let parent_store = self.parent_store.as_ref().unwrap();
390        handle = if new_store_info.encrypted_mutations_object_id == INVALID_OBJECT_ID {
391            let handle = ObjectStore::create_object(
392                parent_store,
393                &mut transaction,
394                HandleOptions { skip_journal_checks: true, ..Default::default() },
395                None,
396            )
397            .await?;
398            let oid = handle.object_id();
399            end_transaction = parent_store
400                .new_transaction(
401                    lock_keys![
402                        LockKey::object(parent_store.store_object_id(), oid),
403                        LockKey::object(
404                            parent_store.store_object_id(),
405                            self.store_info_handle_object_id().unwrap(),
406                        ),
407                    ],
408                    txn_options,
409                )
410                .await?;
411            new_store_info.encrypted_mutations_object_id = oid;
412            parent_store.add_to_graveyard(&mut transaction, oid);
413            parent_store.remove_from_graveyard(&mut end_transaction, oid);
414            handle
415        } else {
416            end_transaction = parent_store
417                .new_transaction(
418                    lock_keys![
419                        LockKey::object(
420                            parent_store.store_object_id(),
421                            new_store_info.encrypted_mutations_object_id,
422                        ),
423                        LockKey::object(
424                            parent_store.store_object_id(),
425                            self.store_info_handle_object_id().unwrap(),
426                        ),
427                    ],
428                    txn_options,
429                )
430                .await?;
431            ObjectStore::open_object(
432                parent_store,
433                new_store_info.encrypted_mutations_object_id,
434                HandleOptions { skip_journal_checks: true, ..Default::default() },
435                None,
436            )
437            .await?
438        };
439        transaction
440            .commit()
441            .await
442            .context("Failed to commit create encrypted mutations transaction")?;
443
444        // Append the encrypted mutations, which need to be read from the journal.
445        // This assumes that the journal has no buffered mutations for this store (see Self::lock).
446        let journaled = filesystem
447            .journal()
448            .read_transactions_for_object(self.store_object_id)
449            .await
450            .context("Failed to read encrypted mutations from journal")?;
451        let mut buffer = handle.allocate_buffer(MAX_ENCRYPTED_MUTATIONS_SIZE).await;
452        let mut writer = buffer.writer();
453        EncryptedMutations::from_replayed_mutations(self.store_object_id, journaled)
454            .serialize_with_version(&mut writer)?;
455        let len = writer.position();
456        handle
457            .txn_write(&mut end_transaction, handle.get_size(), buffer.subslice(..len))
458            .await
459            .context("Failed to write encrypted mutations")?;
460
461        self.write_store_info(&mut end_transaction, &new_store_info)
462            .await
463            .context("Failed to write store info")?;
464
465        let mut total_layer_size = 0;
466        for &oid in &new_store_info.layers {
467            total_layer_size += parent_store.get_file_size(oid).await?;
468        }
469        total_layer_size +=
470            layer_size_from_encrypted_mutations_size(handle.get_size() + len as u64);
471
472        reservation_update =
473            ReservationUpdate::new(tree::reservation_amount_from_layer_size(total_layer_size));
474
475        end_transaction.add_with_object(
476            self.store_object_id(),
477            Mutation::EndFlush,
478            AssocObj::Borrowed(&reservation_update),
479        );
480
481        end_transaction.commit().await.context("Failed to commit EndFlush transaction")?;
482
483        Ok(())
484    }
485}
486
487#[cfg(test)]
488mod tests {
489    use super::{FlushReason, ForceMajor};
490    use crate::filesystem::{FxFilesystem, FxFilesystemBuilder};
491    use crate::object_handle::{
492        INVALID_OBJECT_ID, ObjectHandle, ReadObjectHandle, WriteObjectHandle,
493    };
494    use crate::object_store::directory::Directory;
495    use crate::object_store::journal::JournalOptions;
496    use crate::object_store::transaction::{Options, lock_keys};
497    use crate::object_store::volume::root_volume;
498    use crate::object_store::{
499        HandleOptions, LockKey, NewChildStoreOptions, ObjectStore, StoreOptions,
500        layer_size_from_encrypted_mutations_size, tree,
501    };
502    use fxfs_insecure_crypto::new_insecure_crypt;
503    use std::sync::Arc;
504    use storage_device::DeviceHolder;
505    use storage_device::fake_device::FakeDevice;
506
507    #[fuchsia::test]
508    async fn test_flush_when_locked() {
509        let device = DeviceHolder::new(FakeDevice::new(8192, 1024));
510        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
511        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
512        let crypt = Arc::new(new_insecure_crypt());
513        let store = root_volume
514            .new_volume(
515                "test",
516                NewChildStoreOptions {
517                    options: StoreOptions { crypt: Some(crypt.clone()), ..StoreOptions::default() },
518                    ..NewChildStoreOptions::default()
519                },
520            )
521            .await
522            .expect("new_volume failed");
523        let root_dir =
524            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
525        let mut transaction = fs
526            .root_store()
527            .new_transaction(
528                lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
529                Options::default(),
530            )
531            .await
532            .expect("new_transaction failed");
533        let foo = root_dir
534            .create_child_file(&mut transaction, "foo")
535            .await
536            .expect("create_child_file failed");
537        transaction.commit().await.expect("commit failed");
538
539        // When the volume is first created it will include a new mutations key but we want to test
540        // what happens when the encrypted mutations file doesn't contain a new mutations key, so we
541        // flush here.
542        store.flush().await.expect("flush failed");
543
544        let mut transaction = fs
545            .root_store()
546            .new_transaction(
547                lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
548                Options::default(),
549            )
550            .await
551            .expect("new_transaction failed");
552        let bar = root_dir
553            .create_child_file(&mut transaction, "bar")
554            .await
555            .expect("create_child_file failed");
556        transaction.commit().await.expect("commit failed");
557
558        store.lock().await.expect("lock failed");
559
560        // Flushing the store whilst locked should create an encrypted mutations file.
561        store.flush().await.expect("flush failed");
562
563        // Check the reservation.
564        let info = store.load_store_info().await.unwrap();
565        let parent_store = store.parent_store().unwrap();
566        let mut total_layer_size = 0;
567        for &oid in &info.layers {
568            total_layer_size +=
569                parent_store.get_file_size(oid).await.expect("get_file_size failed");
570        }
571        assert_ne!(info.encrypted_mutations_object_id, INVALID_OBJECT_ID);
572        total_layer_size += layer_size_from_encrypted_mutations_size(
573            parent_store
574                .get_file_size(info.encrypted_mutations_object_id)
575                .await
576                .expect("get_file_size failed"),
577        );
578        assert_eq!(
579            fs.object_manager().reservation(store.store_object_id()),
580            Some(tree::reservation_amount_from_layer_size(total_layer_size))
581        );
582
583        // Unlocking the store should replay that encrypted mutations file.
584        store.unlock(crypt).await.expect("unlock failed");
585
586        ObjectStore::open_object(&store, foo.object_id(), HandleOptions::default(), None)
587            .await
588            .expect("open_object failed");
589
590        ObjectStore::open_object(&store, bar.object_id(), HandleOptions::default(), None)
591            .await
592            .expect("open_object failed");
593
594        fs.close().await.expect("close failed");
595    }
596
597    #[fuchsia::test]
598    async fn test_major_compaction_frees_reservation_and_merges_layers() {
599        let device = DeviceHolder::new(FakeDevice::new(32768, 512));
600        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
601        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
602        let store = root_volume
603            .new_volume("test", NewChildStoreOptions::default())
604            .await
605            .expect("new_volume failed");
606        let root_dir =
607            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
608
609        // 1. Populate the store with files to create an initial layer > 512 KiB
610        // (DEFAULT_RECLAIM_SIZE).
611        let num_files = 2000;
612        let mut file_ids = Vec::with_capacity(num_files);
613        for i in 0..num_files {
614            let filename = format!("file_{:04}_{:<200}", i, i);
615            let mut transaction = store
616                .new_transaction(
617                    lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
618                    Options::default(),
619                )
620                .await
621                .expect("new_transaction failed");
622            let file = root_dir
623                .create_child_file(&mut transaction, &filename)
624                .await
625                .expect("create_child_file failed");
626            transaction.commit().await.expect("commit failed");
627            file_ids.push((file.object_id(), filename));
628        }
629
630        // Flush (minor) to create Layer 0.
631        store.flush().await.expect("flush failed");
632
633        let info_initial = store.load_store_info().await.expect("load_store_info failed");
634        assert_eq!(info_initial.layers.len(), 1);
635        let initial_reservation = fs
636            .object_manager()
637            .reservation(store.store_object_id())
638            .expect("reservation not found");
639
640        // 2. Delete half the files (1000 files).
641        for (oid, filename) in &file_ids[0..1000] {
642            let mut transaction = store
643                .new_transaction(
644                    lock_keys![
645                        LockKey::object(store.store_object_id(), root_dir.object_id()),
646                        LockKey::object(store.store_object_id(), *oid),
647                    ],
648                    Options::default(),
649                )
650                .await
651                .expect("new_transaction failed");
652            let replaced = crate::object_store::directory::replace_child(
653                &mut transaction,
654                None,
655                (&root_dir, filename.as_str()),
656            )
657            .await
658            .expect("replace_child failed");
659            assert_matches::assert_matches!(
660                replaced,
661                crate::object_store::directory::ReplacedChild::Object(id) if id == *oid
662            );
663            transaction.commit().await.expect("commit failed");
664        }
665
666        // 3. Minor flush after deleting 1000 files (FlushReason::Journal(ForceMajor::False)).
667        // Minor compaction will NOT merge with the base layer because base layer > 512 KiB.
668        store
669            .flush_with_reason(FlushReason::Journal(ForceMajor::False))
670            .await
671            .expect("minor flush failed");
672
673        let info_after_minor = store.load_store_info().await.expect("load_store_info failed");
674        assert_eq!(
675            info_after_minor.layers.len(),
676            2,
677            "Minor compaction should keep the base layer, creating 2 layers"
678        );
679        let res_after_minor = fs
680            .object_manager()
681            .reservation(store.store_object_id())
682            .expect("reservation not found");
683        assert!(
684            res_after_minor >= initial_reservation,
685            "Minor compaction should not decrease reservation because base layer is kept"
686        );
687
688        // 4. Major flush (FlushReason::Journal(ForceMajor::True)).
689        // Major compaction must merge ALL layers, purge tombstones, and reduce reservation.
690        store
691            .flush_with_reason(FlushReason::Journal(ForceMajor::True))
692            .await
693            .expect("major flush failed");
694
695        let info_after_major = store.load_store_info().await.expect("load_store_info failed");
696        assert_eq!(
697            info_after_major.layers.len(),
698            1,
699            "Major compaction should merge all layers into 1"
700        );
701        let res_after_major = fs
702            .object_manager()
703            .reservation(store.store_object_id())
704            .expect("reservation not found");
705        assert!(
706            res_after_major < initial_reservation,
707            "Major compaction should decrease reservation (was {}, initial was {})",
708            res_after_major,
709            initial_reservation
710        );
711
712        // 5. Verify data integrity.
713        // Deleted files must be gone.
714        for (oid, filename) in &file_ids[0..1000] {
715            assert!(
716                root_dir.lookup(filename).await.expect("lookup failed").is_none(),
717                "Deleted file {} should not exist",
718                oid
719            );
720        }
721        for (oid, filename) in &file_ids[1000..num_files] {
722            let res = root_dir.lookup(filename).await.expect("lookup failed");
723            assert!(res.is_some(), "Surviving file {} should exist", oid);
724            assert_eq!(res.unwrap().0, *oid);
725        }
726
727        fs.close().await.expect("close failed");
728    }
729
730    #[fuchsia::test]
731    async fn test_major_compaction_without_journal_mutations() {
732        let device = DeviceHolder::new(FakeDevice::new(32768, 512));
733        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
734        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
735        let store = root_volume
736            .new_volume("test", NewChildStoreOptions::default())
737            .await
738            .expect("new_volume failed");
739        let root_dir =
740            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
741
742        // Create initial base layer > 512 KiB.
743        for i in 0..2000 {
744            let filename = format!("file_{:04}_{:<200}", i, i);
745            let mut transaction = store
746                .new_transaction(
747                    lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
748                    Options::default(),
749                )
750                .await
751                .expect("new_transaction failed");
752            root_dir.create_child_file(&mut transaction, &filename).await.expect("create failed");
753            transaction.commit().await.expect("commit failed");
754        }
755        store.flush().await.expect("flush failed");
756
757        // Write a small file and do a minor flush to create a second layer.
758        let mut transaction = store
759            .new_transaction(
760                lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
761                Options::default(),
762            )
763            .await
764            .expect("new_transaction failed");
765        root_dir.create_child_file(&mut transaction, "extra_file").await.expect("create failed");
766        transaction.commit().await.expect("commit failed");
767        store
768            .flush_with_reason(FlushReason::Journal(ForceMajor::False))
769            .await
770            .expect("flush failed");
771
772        let info = store.load_store_info().await.expect("load_store_info failed");
773        assert_eq!(info.layers.len(), 2, "Should have 2 layers before major flush");
774
775        // Now, do a major compaction when there are NO journal mutations in ObjectManager.
776        assert!(!fs.object_manager().needs_flush(store.store_object_id()));
777        store
778            .flush_with_reason(FlushReason::Journal(ForceMajor::True))
779            .await
780            .expect("major flush failed");
781
782        let info_after_major = store.load_store_info().await.expect("load_store_info failed");
783        assert_eq!(
784            info_after_major.layers.len(),
785            1,
786            "Major flush should merge 2 layers into 1 even without journal mutations"
787        );
788
789        // Doing another major flush when already at 1 layer and no mutations should succeed
790        // and leave 1 layer.
791        assert!(!fs.object_manager().needs_flush(store.store_object_id()));
792        store
793            .flush_with_reason(FlushReason::Journal(ForceMajor::True))
794            .await
795            .expect("second major flush failed");
796        let info_after_second = store.load_store_info().await.expect("load_store_info failed");
797        assert_eq!(
798            info_after_second.layers.len(),
799            1,
800            "Major flush when at 1 layer should succeed and leave 1 layer"
801        );
802
803        fs.close().await.expect("close failed");
804    }
805
806    #[fuchsia::test]
807    async fn test_major_compaction_purges_deleted_extents() {
808        let device = DeviceHolder::new(FakeDevice::new(32768, 512));
809        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
810        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
811        let store = root_volume
812            .new_volume("test", NewChildStoreOptions::default())
813            .await
814            .expect("new_volume failed");
815        let root_dir =
816            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
817
818        // 1. Create a base layer > 512 KiB so minor compaction won't merge with it.
819        for i in 0..2000 {
820            let filename = format!("file_{:04}_{:<200}", i, i);
821            let mut transaction = store
822                .new_transaction(
823                    lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
824                    Options::default(),
825                )
826                .await
827                .expect("new_transaction failed");
828            root_dir.create_child_file(&mut transaction, &filename).await.expect("create failed");
829            transaction.commit().await.expect("commit failed");
830        }
831
832        // Create a data file with extents.
833        let data_oid = {
834            let mut transaction = store
835                .new_transaction(
836                    lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
837                    Options::default(),
838                )
839                .await
840                .expect("new_transaction failed");
841            let file = root_dir
842                .create_child_file(&mut transaction, "data_file")
843                .await
844                .expect("create_child_file failed");
845            let oid = file.object_id();
846            transaction.commit().await.expect("commit failed");
847            oid
848        };
849
850        let data_file = ObjectStore::open_object(&store, data_oid, HandleOptions::default(), None)
851            .await
852            .expect("open_object failed");
853        {
854            let mut transaction = store
855                .new_transaction(
856                    lock_keys![LockKey::object(store.store_object_id(), data_oid)],
857                    Options::default(),
858                )
859                .await
860                .expect("new_transaction failed");
861            let mut buffer = data_file.allocate_buffer(131072).await;
862            buffer.fill(0xAB);
863            data_file
864                .txn_write(&mut transaction, 0, buffer.as_ref())
865                .await
866                .expect("txn_write failed");
867            transaction.commit().await.expect("commit failed");
868        }
869
870        // Flush (minor) to commit extents to Layer 0.
871        store.flush().await.expect("flush failed");
872
873        let info0 = store.load_store_info().await.expect("load_store_info failed");
874        assert_eq!(info0.layers.len(), 1);
875
876        // Truncate data_file to 0 to write extent tombstones (ExtentValue::None).
877        data_file.truncate(0).await.expect("truncate failed");
878
879        // Minor flush: writes Layer 1 with extent tombstones on top of Layer 0.
880        store
881            .flush_with_reason(FlushReason::Journal(ForceMajor::False))
882            .await
883            .expect("minor flush failed");
884        let info1 = store.load_store_info().await.expect("load_store_info failed");
885        assert_eq!(info1.layers.len(), 2, "Minor flush should keep Layer 0 and create Layer 1");
886
887        // Major flush: purges ExtentValue::None tombstones via major_iter and merges into 1 layer.
888        store
889            .flush_with_reason(FlushReason::Journal(ForceMajor::True))
890            .await
891            .expect("major flush failed");
892        let info2 = store.load_store_info().await.expect("load_store_info failed");
893        assert_eq!(info2.layers.len(), 1, "Major flush should merge into 1 layer");
894
895        // Verify data_file size is 0.
896        assert_eq!(data_file.get_size(), 0);
897
898        fs.close().await.expect("close failed");
899    }
900
901    // This test case is a regression test for b/548631578.  It verifies that we will force major
902    // compactions when borrowed_metadata_space gets too large relative to metadata_reservation.
903    // Minor compactions are not sufficient in all cases to return borrowed metadata space, which
904    // can eventually exhuast the metadata reservation and prevent all operations (including
905    // compaction itself).
906    #[fuchsia::test]
907    async fn test_borrow_metadata_space_fails_without_major_compaction() {
908        let reclaim_size = 65536;
909        // 102,400 blocks of 512 bytes = 50 MiB device.
910        let device = DeviceHolder::new(FakeDevice::new(102400, 512));
911        let fs = FxFilesystemBuilder::new()
912            .journal_options(JournalOptions { reclaim_size, ..Default::default() })
913            .format(true)
914            .open(device)
915            .await
916            .expect("open failed");
917
918        // The test creates many files in several stores, and then deletes them.  The deletions
919        // create additional layer files, which should be collapsed into the base layers (and
920        // cancel out the file creations) following a major compaction, which should happen
921        // automatically due to the large amount of borrowed space.
922        let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
923        let num_stores = 3;
924        let files_per_store = 3500;
925        let mut stores_and_files = Vec::new();
926
927        for s in 0..num_stores {
928            let store = root_volume
929                .new_volume(&format!("test_{s}"), NewChildStoreOptions::default())
930                .await
931                .expect("new_volume failed");
932            let root_dir = Directory::open(&store, store.root_directory_object_id())
933                .await
934                .expect("open failed");
935            let mut file_ids = Vec::with_capacity(files_per_store);
936            for i in 0..files_per_store {
937                let filename = format!("file_{:04}_{:<200}", i, i);
938                let mut transaction = store
939                    .new_transaction(
940                        lock_keys![LockKey::object(store.store_object_id(), root_dir.object_id())],
941                        Options::default(),
942                    )
943                    .await
944                    .expect("new_transaction failed");
945                let file = root_dir
946                    .create_child_file(&mut transaction, &filename)
947                    .await
948                    .expect("create_child_file failed");
949                transaction.commit().await.expect("commit failed");
950                file_ids.push((file.object_id(), filename));
951            }
952            store
953                .flush_with_reason(FlushReason::Journal(ForceMajor::True))
954                .await
955                .expect("flush failed");
956            stores_and_files.push((store, root_dir, file_ids));
957        }
958
959        for (_, root_dir, file_ids) in stores_and_files {
960            for (oid, filename) in &file_ids[0..2500] {
961                let mut context = root_dir
962                    .acquire_context_for_replace(None, filename.as_str(), true)
963                    .await
964                    .expect("acquire_context failed");
965                let replaced = crate::object_store::directory::replace_child(
966                    &mut context.transaction,
967                    None,
968                    (&root_dir, filename.as_str()),
969                )
970                .await
971                .expect("replace_child failed");
972                assert_matches::assert_matches!(
973                    replaced,
974                    crate::object_store::directory::ReplacedChild::Object(id) if id == *oid
975                );
976                context.transaction.commit().await.expect("commit failed");
977            }
978        }
979
980        fs.close().await.expect("close failed");
981    }
982}
983
984impl tree::MajorCompactable<ObjectKey, ObjectValue> for LSMTree<ObjectKey, ObjectValue> {
985    async fn major_iter(
986        iter: impl LayerIterator<ObjectKey, ObjectValue>,
987    ) -> Result<impl LayerIterator<ObjectKey, ObjectValue>, Error> {
988        iter.filter(|item: ItemRef<'_, _, _>| match item {
989            // Object Tombstone.
990            ItemRef { value: ObjectValue::None, .. } => false,
991            // Deleted extent.
992            ItemRef { value: ObjectValue::Extent(ExtentValue::None), .. } => false,
993            _ => true,
994        })
995        .await
996    }
997}