Skip to main content

fxfs/
filesystem.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
5use crate::errors::FxfsError;
6use crate::fsck::{FsckOptions, fsck_volume_with_options, fsck_with_options};
7use crate::hooks::HooksHandle;
8use crate::log::*;
9use crate::metrics;
10use crate::object_handle::LayerObject;
11use crate::object_store::allocator::Allocator;
12use crate::object_store::directory::Directory;
13use crate::object_store::graveyard::Graveyard;
14use crate::object_store::journal::super_block::{SuperBlockHeader, SuperBlockInstance};
15use crate::object_store::journal::{self, Journal, JournalCheckpoint, JournalOptions};
16use crate::object_store::object_manager::ObjectManager;
17use crate::object_store::transaction::{
18    self, AssocObj, LockKey, LockManager, MetadataReservation, Mutation, ObjectMutationIterator,
19    TRANSACTION_METADATA_MAX_AMOUNT, Transaction, WriteGuard, lock_keys,
20};
21use crate::object_store::volume::{VOLUMES_DIRECTORY, root_volume};
22use crate::object_store::{
23    AttributeId, DataObjectHandle, NewChildStoreOptions, ObjectStore, StoreOptions,
24};
25use crate::range::RangeExt;
26use crate::serialized_types::{LATEST_VERSION, Version};
27use anyhow::{Context, Error, anyhow, bail};
28use async_trait::async_trait;
29use event_listener::Event;
30use fuchsia_async as fasync;
31use fuchsia_async::condition::Condition;
32use fuchsia_inspect::{Inspector, LazyNode, NumericProperty as _, UintProperty};
33use fuchsia_sync::Mutex;
34use futures::future::BoxFuture;
35use futures::stream::BoxStream;
36use futures::{FutureExt, Stream};
37use fxfs_crypto::{Crypt, UnwrappedKey};
38use fxfs_trace::{TraceFutureExt, trace_future_args};
39use static_assertions::const_assert;
40use std::pin::pin;
41use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
42use std::sync::{Arc, OnceLock, Weak};
43use std::task::Poll;
44use std::time::{Duration, Instant};
45use storage_device::{Device, DeviceHolder};
46use storage_units::BlockSize;
47
48pub const MIN_BLOCK_SIZE: BlockSize = BlockSize::SIZE_4KIB;
49pub const MAX_BLOCK_SIZE: BlockSize = BlockSize::SIZE_64KIB;
50
51// Whilst Fxfs could support up to u64::MAX, off_t is i64 so allowing files larger than that becomes
52// difficult to deal with via the POSIX APIs. Additionally, PagedObjectHandle only sees data get
53// modified in page chunks so to prevent writes at i64::MAX the entire page containing i64::MAX
54// needs to be excluded.
55pub const MAX_FILE_SIZE: u64 = i64::MAX as u64 - 4095;
56const_assert!(9223372036854771712 == MAX_FILE_SIZE);
57
58use futures::stream::StreamExt;
59
60// The maximum number of transactions that can be in-flight at any time.
61pub(crate) const MAX_IN_FLIGHT_TRANSACTIONS: u64 = 4;
62
63// Start trimming 1 hour after boot.  The idea here is to wait until the initial flurry of
64// activity during boot is finished.  This is a rough heuristic and may need to change later if
65// performance is affected.
66const TRIM_AFTER_BOOT_TIMER: Duration = Duration::from_secs(60 * 60);
67
68// After the initial trim, perform another trim every 24 hours.
69const TRIM_INTERVAL_TIMER: Duration = Duration::from_secs(60 * 60 * 24);
70
71/// How often to clean the transfer buffer.
72// TODO(https://fxbug.dev/489725256) Configure the task to run when fxfs is idle.
73const CLEAN_TRANSFER_BUFFER_INTERVAL: Duration = Duration::from_secs(60);
74
75/// An opaque token representing an active wake lease. When held, it prevents the system from
76/// suspending. Dropping this token releases the lease.
77#[derive(Debug)]
78pub struct WakeLease {
79    #[cfg(target_os = "fuchsia")]
80    _lease: zx::Handle,
81}
82
83#[cfg(target_os = "fuchsia")]
84impl WakeLease {
85    pub fn new<H: TryInto<zx::Handle>>(lease: H) -> Self
86    where
87        H::Error: std::fmt::Debug,
88    {
89        Self { _lease: lease.try_into().expect("lease handle must be valid") }
90    }
91}
92
93pub trait PowerManager: Send + Sync {
94    /// Returns a stream of battery status changes (true if using battery).
95    fn watch_battery(self: Arc<Self>) -> BoxStream<'static, (bool, Option<WakeLease>)>;
96}
97
98/// Services registration of layer files with a platform-specific layer pager.
99#[async_trait]
100pub trait LayerPager: Send + Sync + 'static {
101    async fn open_layer(
102        &self,
103        handle: DataObjectHandle<ObjectStore>,
104        unwrapped_key: Option<UnwrappedKey>,
105    ) -> Result<Arc<dyn LayerObject>, Error>;
106}
107
108/// Holds information on an Fxfs Filesystem
109pub struct Info {
110    pub total_bytes: u64,
111    pub used_bytes: u64,
112}
113
114pub type PostCommitHook = Option<Box<dyn Fn() -> BoxFuture<'static, ()> + Send + Sync>>;
115
116pub struct Options {
117    /// True if the filesystem is read-only.
118    pub read_only: bool,
119
120    /// Hooks for filesystem events (e.g. pre_commit, before_commit).
121    pub hooks: Arc<HooksHandle>,
122
123    /// A callback that runs after every transaction has been committed.  This will be called whilst
124    /// a lock is held which will block more transactions from being committed.
125    pub post_commit_hook: PostCommitHook,
126
127    /// If true, don't do an initial reap of the graveyard at mount time.  This is useful for
128    /// testing.
129    pub skip_initial_reap: bool,
130
131    // The first duration is how long after the filesystem has been mounted to perform an initial
132    // trim.  The second is the interval to repeat trimming thereafter.  If set to None, no trimming
133    // is done.
134    // Default values are (5 minutes, 24 hours).
135    pub trim_config: Option<(Duration, Duration)>,
136
137    // If set, journal will not be used for writes. The user must call 'close' when finished.
138    // The provided superblock instance will be written upon close().
139    pub image_builder_mode: Option<SuperBlockInstance>,
140
141    // If true, the filesystem will use the hardware's inline crypto engine to write encrypted
142    // data. Requires the block device to support inline encryption and for `barriers_enabled` to
143    // be true.
144    // TODO(https://fxbug.dev/393196849): For now, this flag only prevents the filesystem from
145    // computing checksums. Update this comment when the filesystem actually uses inline
146    // encryption.
147    pub inline_crypto_enabled: bool,
148
149    // Configures the filesystem to use barriers instead of checksums to ensure consistency.
150    // Checksums may be computed and stored in extent records but will no longer be stored in the
151    // journal. The journal will use barriers to enforce proper ordering between data and metadata
152    // writes. Must be true if `inline_crypto_enabled` is true.
153    pub barriers_enabled: bool,
154
155    /// If set, this will be used to check for charger status before trimming.
156    pub power_manager: Option<Arc<dyn PowerManager>>,
157
158    /// How long to wait after being placed on a charger before starting a trim.
159    pub trim_charger_wait: Duration,
160
161    /// If true, allows writing Type 3 delivery blobs.
162    /// NOTE: Type 3 delivery blobs are currently UNSTABLE / EXPERIMENTAL and subject to change.
163    pub allow_type3_blobs: bool,
164}
165
166impl Default for Options {
167    fn default() -> Self {
168        Options {
169            read_only: false,
170            hooks: Arc::<HooksHandle>::default(),
171            post_commit_hook: None,
172            skip_initial_reap: false,
173            trim_config: Some((TRIM_AFTER_BOOT_TIMER, TRIM_INTERVAL_TIMER)),
174            image_builder_mode: None,
175            inline_crypto_enabled: false,
176            barriers_enabled: false,
177            power_manager: None,
178            trim_charger_wait: Duration::from_secs(10),
179            allow_type3_blobs: false,
180        }
181    }
182}
183
184/// The context in which a transaction is being applied.
185pub struct ApplyContext<'a, 'b> {
186    /// The mode indicates whether the transaction is being replayed.
187    pub mode: ApplyMode<'a, 'b>,
188
189    /// The transaction checkpoint for this mutation.
190    pub checkpoint: JournalCheckpoint,
191}
192
193/// A transaction can be applied during replay or on a live running system (in which case a
194/// transaction object will be available).
195pub enum ApplyMode<'a, 'b> {
196    Replay,
197    Live(&'a Transaction<'b>),
198}
199
200impl ApplyMode<'_, '_> {
201    pub fn is_replay(&self) -> bool {
202        matches!(self, ApplyMode::Replay)
203    }
204
205    pub fn is_live(&self) -> bool {
206        matches!(self, ApplyMode::Live(_))
207    }
208}
209
210#[derive(Copy, Clone, Debug, PartialEq, Eq)]
211pub enum ForceMajor {
212    True,
213    False,
214}
215
216#[derive(Copy, Clone, Debug, PartialEq, Eq)]
217pub enum FlushReason {
218    /// Journal memory or space pressure.
219    Journal(ForceMajor),
220
221    /// Clean up an encrypted mutations object after mount.
222    EncryptedMutations,
223
224    /// Upgrade old layer files to the latest version after mount. This performs a full compaction.
225    UpgradeVersion,
226}
227
228/// Objects that use journaling to track mutations (`Allocator` and `ObjectStore`) implement this.
229/// This is primarily used by `ObjectManager` and `SuperBlock` with flush calls used in a few tests.
230#[async_trait]
231pub trait JournalingObject: Send + Sync {
232    /// This method get called when the transaction commits, which can either be during live
233    /// operation (See `ObjectManager::apply_mutation`) or during journal replay, in which case
234    /// transaction will be None (See `super_block::read`).
235    fn apply_mutation(
236        &self,
237        mutation: Mutation,
238        context: &ApplyContext<'_, '_>,
239        assoc_obj: AssocObj<'_>,
240    ) -> Result<(), Error>;
241
242    /// Called when a transaction fails to commit.
243    fn drop_mutation(&self, mutation: Mutation, transaction: &Transaction<'_>);
244
245    /// Called before committing a transaction. Implementations can use this to acquire locks
246    /// or resources (like keys) that must be held until the transaction is committed.
247    async fn prepare_commit<'a>(
248        &self,
249        _filesystem: &'a FxFilesystem,
250        _transaction: &Transaction<'_>,
251    ) -> Result<Option<WriteGuard<'a>>, Error> {
252        Ok(None)
253    }
254
255    /// Flushes in-memory changes to the device (to allow journal space to be freed).
256    ///
257    /// Also returns the earliest version of a struct in the filesystem.
258    async fn flush(&self, reason: FlushReason) -> Result<Version, Error>;
259
260    /// Writes mutations to the journal.  This allows objects to encrypt or otherwise modify what
261    /// gets written to the journal.
262    fn write_mutations(
263        &self,
264        mutations: ObjectMutationIterator<'_, '_>,
265        mut writer: journal::Writer<'_>,
266    ) {
267        for mutation in mutations {
268            writer.write(mutation.clone());
269        }
270    }
271}
272
273#[derive(Default)]
274pub struct SyncOptions<'a> {
275    /// If set, the journal will be flushed, as well as the underlying block device.  This is much
276    /// more expensive, but ensures the contents of the journal are persisted (which also acts as a
277    /// barrier, ensuring all previous journal writes are observable by future operations).
278    /// Note that when this is not set, the journal is *not* synchronously flushed by the sync call,
279    /// and it will return before the journal flush completes.  In other words, some journal
280    /// mutations may still be buffered in memory after this call returns.
281    pub flush_device: bool,
282
283    /// A precondition that is evaluated whilst a lock is held that determines whether or not the
284    /// sync needs to proceed.
285    pub precondition: Option<Box<dyn FnOnce() -> bool + 'a + Send>>,
286}
287
288pub struct OpenFxFilesystem(Arc<FxFilesystem>);
289
290impl OpenFxFilesystem {
291    /// Waits for filesystem to be dropped (so callers should ensure all direct and indirect
292    /// references are dropped) and returns the device.  No attempt is made at a graceful shutdown.
293    pub async fn take_device(self) -> DeviceHolder {
294        let fut = self.device.take_when_dropped();
295        std::mem::drop(self);
296        debug_assert_not_too_long!(fut)
297    }
298}
299
300impl From<Arc<FxFilesystem>> for OpenFxFilesystem {
301    fn from(fs: Arc<FxFilesystem>) -> Self {
302        Self(fs)
303    }
304}
305
306impl Drop for OpenFxFilesystem {
307    fn drop(&mut self) {
308        if self.options.image_builder_mode.is_some()
309            && self.journal().image_builder_mode().is_some()
310        {
311            error!("OpenFxFilesystem in image_builder_mode dropped without calling close().");
312        }
313        if !self.options.read_only && !self.closed.load(Ordering::SeqCst) {
314            error!("OpenFxFilesystem dropped without first being closed. Data loss may occur.");
315        }
316    }
317}
318
319impl std::ops::Deref for OpenFxFilesystem {
320    type Target = Arc<FxFilesystem>;
321
322    fn deref(&self) -> &Self::Target {
323        &self.0
324    }
325}
326
327pub struct FxFilesystemBuilder {
328    format: bool,
329    trace: bool,
330    options: Options,
331    journal_options: JournalOptions,
332    on_new_allocator: Option<Box<dyn Fn(Arc<Allocator>) + Send + Sync>>,
333    on_new_store: Option<Box<dyn Fn(&ObjectStore) + Send + Sync>>,
334    fsck_after_every_transaction: bool,
335    layer_pager: Option<Arc<dyn LayerPager>>,
336}
337
338impl FxFilesystemBuilder {
339    pub fn new() -> Self {
340        Self {
341            format: false,
342            trace: false,
343            options: Options::default(),
344            journal_options: JournalOptions::default(),
345            on_new_allocator: None,
346            on_new_store: None,
347            fsck_after_every_transaction: false,
348            layer_pager: None,
349        }
350    }
351
352    /// Sets whether the block device should be formatted when opened. Defaults to `false`.
353    pub fn format(mut self, format: bool) -> Self {
354        self.format = format;
355        self
356    }
357
358    /// Enables or disables trace level logging. Defaults to `false`.
359    pub fn trace(mut self, trace: bool) -> Self {
360        self.trace = trace;
361        self
362    }
363
364    /// Sets whether the filesystem will be opened in read-only mode. Defaults to `false`.
365    /// Incompatible with `format`.
366    pub fn read_only(mut self, read_only: bool) -> Self {
367        self.options.read_only = read_only;
368        self
369    }
370
371    /// Sets whether Type 3 delivery blobs are allowed. Defaults to `false`.
372    pub fn allow_type3_blobs(mut self, allow: bool) -> Self {
373        self.options.allow_type3_blobs = allow;
374        self
375    }
376
377    /// For image building and in-place migration.
378    ///
379    /// This mode avoids the initial write of super blocks and skips the journal for all
380    /// transactions. The user *must* call `close()` before dropping the filesystem to trigger
381    /// a compaction of in-memory data structures, a minimal journal and a write to one
382    /// superblock (as specified).
383    pub fn image_builder_mode(mut self, mode: Option<SuperBlockInstance>) -> Self {
384        self.options.image_builder_mode = mode;
385        self
386    }
387
388    /// Sets a callback that runs after every transaction has been committed. See
389    /// `Options::post_commit_hook`.
390    pub fn post_commit_hook(
391        mut self,
392        hook: impl Fn() -> futures::future::BoxFuture<'static, ()> + Send + Sync + 'static,
393    ) -> Self {
394        self.options.post_commit_hook = Some(Box::new(hook));
395        self
396    }
397
398    /// Sets whether to do an initial reap of the graveyard at mount time. See
399    /// `Options::skip_initial_reap`. Defaults to `false`.
400    pub fn skip_initial_reap(mut self, skip_initial_reap: bool) -> Self {
401        self.options.skip_initial_reap = skip_initial_reap;
402        self
403    }
404
405    /// Sets the options for the journal.
406    pub fn journal_options(mut self, journal_options: JournalOptions) -> Self {
407        self.journal_options = journal_options;
408        self
409    }
410
411    /// Sets a method to be called immediately after creating the allocator.
412    pub fn on_new_allocator(
413        mut self,
414        on_new_allocator: impl Fn(Arc<Allocator>) + Send + Sync + 'static,
415    ) -> Self {
416        self.on_new_allocator = Some(Box::new(on_new_allocator));
417        self
418    }
419
420    /// Sets a method to be called each time a new store is registered with `ObjectManager`.
421    pub fn on_new_store(
422        mut self,
423        on_new_store: impl Fn(&ObjectStore) + Send + Sync + 'static,
424    ) -> Self {
425        self.on_new_store = Some(Box::new(on_new_store));
426        self
427    }
428
429    /// Enables or disables running fsck after every transaction. Defaults to `false`.
430    pub fn fsck_after_every_transaction(mut self, fsck_after_every_transaction: bool) -> Self {
431        self.fsck_after_every_transaction = fsck_after_every_transaction;
432        self
433    }
434
435    pub fn trim_config(mut self, delay_and_interval: Option<(Duration, Duration)>) -> Self {
436        self.options.trim_config = delay_and_interval;
437        self
438    }
439
440    pub fn power_manager(mut self, power_manager: Arc<dyn PowerManager>) -> Self {
441        self.options.power_manager = Some(power_manager);
442        self
443    }
444
445    pub fn trim_charger_wait(mut self, wait: Duration) -> Self {
446        self.options.trim_charger_wait = wait;
447        self
448    }
449
450    /// Enables or disables inline encryption. Defaults to `false`.
451    pub fn inline_crypto_enabled(mut self, inline_crypto_enabled: bool) -> Self {
452        self.options.inline_crypto_enabled = inline_crypto_enabled;
453        self
454    }
455
456    /// Enables or disables barriers in both the filesystem and journal options.
457    /// Defaults to `false`.
458    pub fn barriers_enabled(mut self, barriers_enabled: bool) -> Self {
459        self.options.barriers_enabled = barriers_enabled;
460        self.journal_options.barriers_enabled = barriers_enabled;
461        self
462    }
463
464    pub fn hooks(mut self, hooks: Arc<crate::hooks::HooksHandle>) -> Self {
465        self.options.hooks = hooks;
466        self
467    }
468
469    pub fn layer_pager(mut self, layer_pager: Option<Arc<dyn LayerPager>>) -> Self {
470        self.layer_pager = layer_pager;
471        self
472    }
473
474    /// Constructs an `FxFilesystem` object with the specified settings.
475    pub async fn open(self, device: DeviceHolder) -> Result<OpenFxFilesystem, Error> {
476        let read_only = self.options.read_only;
477        if self.format && read_only {
478            bail!("Cannot initialize a filesystem as read-only");
479        }
480
481        // Inline encryption requires barriers to be enabled.
482        if self.options.inline_crypto_enabled && !self.options.barriers_enabled {
483            bail!("A filesystem using inline encryption requires barriers");
484        }
485
486        let objects = Arc::new(ObjectManager::new(self.on_new_store));
487        let journal = Arc::new(Journal::new(
488            objects.clone(),
489            self.journal_options,
490            self.options.hooks.clone(),
491        ));
492
493        let image_builder_mode = self.options.image_builder_mode;
494
495        let device_block_size =
496            BlockSize::new(device.block_size()).expect("Device block size is not a power of 2");
497        let block_size = std::cmp::max(device_block_size, MIN_BLOCK_SIZE);
498        assert!(block_size <= MAX_BLOCK_SIZE, "Max supported block size is 64KiB");
499
500        let mut fsck_after_every_transaction = None;
501        let mut filesystem_options = self.options;
502        if self.fsck_after_every_transaction {
503            let instance =
504                FsckAfterEveryTransaction::new(filesystem_options.post_commit_hook.take());
505            fsck_after_every_transaction = Some(instance.clone());
506            filesystem_options.post_commit_hook =
507                Some(Box::new(move || Box::pin(instance.clone().run())));
508        }
509
510        if !read_only && !self.format {
511            // See comment in JournalRecord::DidFlushDevice for why we need to flush the device
512            // before replay.
513            device.flush().await.context("Device flush failed")?;
514        }
515
516        let filesystem = Arc::new_cyclic(|weak: &Weak<FxFilesystem>| {
517            let weak = weak.clone();
518            FxFilesystem {
519                device,
520                block_size,
521                objects: objects.clone(),
522                journal,
523                commit_mutex: futures::lock::Mutex::new(()),
524                lock_manager: LockManager::new(),
525                flush_task: Mutex::new(None),
526                background_tasks: fasync::Scope::new(),
527                closed: AtomicBool::new(true),
528                trace: self.trace,
529                graveyard: Graveyard::new(weak.clone()),
530                completed_transactions: metrics::detail().create_uint("completed_transactions", 0),
531                options: filesystem_options,
532                in_flight_transactions: AtomicU64::new(0),
533                transaction_limit_event: Event::new(),
534                layer_pager: self.layer_pager,
535                _stores_node: metrics::register_fs(move || {
536                    let weak = weak.clone();
537                    Box::pin(async move {
538                        if let Some(fs) = weak.upgrade() {
539                            fs.populate_stores_node().await
540                        } else {
541                            Err(anyhow!("Filesystem has been dropped"))
542                        }
543                    })
544                }),
545            }
546        });
547
548        filesystem.journal().set_image_builder_mode(image_builder_mode);
549
550        filesystem.journal.set_trace(self.trace);
551        if self.format {
552            filesystem.journal.init_empty(filesystem.clone()).await?;
553            if image_builder_mode.is_none() {
554                // The filesystem isn't valid until superblocks are written but we want to defer
555                // that until last when migrating filesystems or building system images.
556                filesystem.journal.init_superblocks().await?;
557
558                // Start the graveyard's background reaping task.
559                filesystem.graveyard.clone().reap_async();
560            }
561
562            // Create the root volume directory.
563            let root_store = filesystem.root_store();
564            root_store.set_trace(self.trace);
565            let root_directory =
566                Directory::open(&root_store, root_store.root_directory_object_id())
567                    .await
568                    .context("Unable to open root volume directory")?;
569            let mut transaction = root_store
570                .new_transaction(
571                    lock_keys![LockKey::object(
572                        root_store.store_object_id(),
573                        root_directory.object_id()
574                    )],
575                    transaction::Options::default(),
576                )
577                .await?;
578            let volume_directory =
579                root_directory.create_child_dir(&mut transaction, VOLUMES_DIRECTORY).await?;
580            transaction.commit().await?;
581            objects.set_volume_directory(volume_directory);
582        } else {
583            filesystem
584                .journal
585                .replay(filesystem.clone(), self.on_new_allocator)
586                .await
587                .context("Journal replay failed")?;
588            filesystem.root_store().set_trace(self.trace);
589
590            if !read_only {
591                // Queue all purged entries for tombstoning.  Don't start the reaper yet because
592                // that can trigger a flush which can add more entries to the graveyard which might
593                // get caught in the initial reap and cause objects to be prematurely tombstoned.
594                for store in objects.unlocked_stores() {
595                    filesystem.graveyard.initial_reap(&store).await?;
596                }
597            }
598        }
599
600        // This must be after we've formatted the filesystem; it will fail during format otherwise.
601        if let Some(fsck_after_every_transaction) = fsck_after_every_transaction {
602            fsck_after_every_transaction
603                .fs
604                .set(Arc::downgrade(&filesystem))
605                .unwrap_or_else(|_| unreachable!());
606        }
607
608        filesystem.closed.store(false, Ordering::SeqCst);
609
610        if !read_only && image_builder_mode.is_none() {
611            // Start the background tasks.
612            filesystem.graveyard.clone().reap_async();
613
614            if filesystem.options.trim_config.is_some() {
615                filesystem.start_trim_task();
616            }
617            filesystem.start_clean_transfer_buffer_task();
618        }
619
620        Ok(filesystem.into())
621    }
622}
623
624pub struct FxFilesystem {
625    block_size: BlockSize,
626    objects: Arc<ObjectManager>,
627    journal: Arc<Journal>,
628    commit_mutex: futures::lock::Mutex<()>,
629    lock_manager: LockManager,
630    flush_task: Mutex<Option<fasync::Task<()>>>,
631    background_tasks: fasync::Scope,
632    closed: AtomicBool,
633    // An event that is signalled when the filesystem starts to shut down.
634    trace: bool,
635    graveyard: Arc<Graveyard>,
636    completed_transactions: UintProperty,
637    options: Options,
638
639    // The number of in-flight transactions which we will limit to MAX_IN_FLIGHT_TRANSACTIONS.
640    in_flight_transactions: AtomicU64,
641
642    // An event that is used to wake up tasks that are blocked due to the in-flight transaction
643    // limit.
644    transaction_limit_event: Event,
645
646    // The "stores" node in the Inspect tree.
647    _stores_node: LazyNode,
648
649    layer_pager: Option<Arc<dyn LayerPager>>,
650
651    // NOTE: This *must* go last so that when users take the device from a closed filesystem, the
652    // filesystem has dropped all other members first (Rust drops members in declaration order).
653    device: DeviceHolder,
654}
655
656#[fxfs_trace::trace]
657impl FxFilesystem {
658    pub async fn new_empty(device: DeviceHolder) -> Result<OpenFxFilesystem, Error> {
659        FxFilesystemBuilder::new().format(true).open(device).await
660    }
661
662    pub async fn open(device: DeviceHolder) -> Result<OpenFxFilesystem, Error> {
663        FxFilesystemBuilder::new().open(device).await
664    }
665
666    pub fn root_parent_store(&self) -> Arc<ObjectStore> {
667        self.objects.root_parent_store()
668    }
669
670    pub fn layer_pager(&self) -> Option<&Arc<dyn LayerPager>> {
671        self.layer_pager.as_ref()
672    }
673
674    pub async fn close(&self) -> Result<(), Error> {
675        if self.journal().image_builder_mode().is_some() {
676            self.journal().allocate_journal().await?;
677            self.journal().set_image_builder_mode(None);
678            self.journal().force_compact().await?;
679        }
680        assert_eq!(self.closed.swap(true, Ordering::SeqCst), false);
681        debug_assert_not_too_long!(self.graveyard.wait_for_reap());
682        debug_assert_not_too_long!(self.background_tasks.clone().cancel());
683        self.journal.stop_compactions().await;
684        let sync_status =
685            if self.journal().image_builder_mode().is_some() || self.options().read_only {
686                Ok(None)
687            } else {
688                self.journal.sync(SyncOptions { flush_device: true, ..Default::default() }).await
689            };
690        match &sync_status {
691            Ok(None) => {}
692            Ok(checkpoint) => info!(
693                "Filesystem closed (checkpoint={}, metadata_reservation={:?}, \
694                 reservation_required={}, borrowed={})",
695                checkpoint.as_ref().unwrap().0.file_offset,
696                self.object_manager().metadata_reservation(),
697                self.object_manager().required_reservation(),
698                self.object_manager().borrowed_metadata_space(),
699            ),
700            Err(e) => error!(error:? = e; "Failed to sync filesystem; data may be lost"),
701        }
702        self.journal.terminate();
703        let flush_task = self.flush_task.lock().take();
704        if let Some(task) = flush_task {
705            debug_assert_not_too_long!(task);
706        }
707        // Regardless of whether sync succeeds, we should close the device, since otherwise we will
708        // crash instead of exiting gracefully.
709        self.device().close().await.context("Failed to close device")?;
710        sync_status.map(|_| ())
711    }
712
713    pub fn device(&self) -> Arc<dyn Device> {
714        Arc::clone(&self.device)
715    }
716
717    pub fn root_store(&self) -> Arc<ObjectStore> {
718        self.objects.root_store()
719    }
720
721    pub fn allocator(&self) -> Arc<Allocator> {
722        self.objects.allocator()
723    }
724
725    /// Enables allocations for the allocator.
726    /// This is only used in image_builder_mode where it *must*
727    /// be called before any allocations can take place.
728    pub fn enable_allocations(&self) {
729        self.allocator().enable_allocations();
730    }
731
732    pub fn object_manager(&self) -> &Arc<ObjectManager> {
733        &self.objects
734    }
735
736    pub fn journal(&self) -> &Arc<Journal> {
737        &self.journal
738    }
739
740    pub async fn sync(&self, options: SyncOptions<'_>) -> Result<(), Error> {
741        self.journal.sync(options).await.map(|_| ())
742    }
743
744    pub fn block_size(&self) -> BlockSize {
745        self.block_size
746    }
747
748    pub fn get_info(&self) -> Info {
749        Info {
750            total_bytes: self.device.size(),
751            used_bytes: self.object_manager().allocator().get_used_bytes().0,
752        }
753    }
754
755    pub fn super_block_header(&self) -> SuperBlockHeader {
756        self.journal.super_block_header()
757    }
758
759    pub fn graveyard(&self) -> &Arc<Graveyard> {
760        &self.graveyard
761    }
762
763    pub fn trace(&self) -> bool {
764        self.trace
765    }
766
767    pub fn options(&self) -> &Options {
768        &self.options
769    }
770
771    pub fn scope(&self) -> &fasync::Scope {
772        &self.background_tasks
773    }
774
775    /// Returns a guard that must be taken before any transaction can commence.  This guard takes a
776    /// shared lock on the filesystem.  `fsck` will take an exclusive lock so that it can get a
777    /// consistent picture of the filesystem that it can verify.  It is important that this lock is
778    /// acquired before *all* other locks.  It is also important that this lock is not taken twice
779    /// by the same task since that can lead to deadlocks if another task tries to take a write
780    /// lock.
781    pub async fn lock_commits(&self) -> futures::lock::MutexGuard<'_, ()> {
782        self.commit_mutex.lock().await
783    }
784
785    #[trace]
786    pub async fn commit_transaction<R: Send>(
787        &self,
788        transaction: &mut Transaction<'_>,
789        callback: impl FnOnce(u64) -> R + Send,
790    ) -> Result<R, Error> {
791        self.hooks().on_pre_commit(transaction)?;
792        debug_assert_not_too_long!(self.lock_manager.commit_prepare(&transaction));
793
794        // Call prepare_commit on all unique objects involved in the transaction.
795        // We must hold the returned guards until the transaction is committed.
796        // Since transaction.mutations() is sorted by object_id, we can deduplicate
797        // on-the-fly.
798        let mut guards = Vec::new();
799        let mut last_object_id = 0;
800        for mutation in transaction.mutations() {
801            let object_id = mutation.object_id;
802
803            // We don't need to prepare commits (which reserves keys) for flush mutations.
804            if matches!(mutation.mutation, Mutation::BeginFlush | Mutation::EndFlush) {
805                continue;
806            }
807
808            if object_id == last_object_id {
809                continue;
810            }
811            assert!(object_id > last_object_id);
812            last_object_id = object_id;
813
814            if let Some(obj) = self.object_manager().journaling_object(object_id) {
815                if let Some(guard) = obj.prepare_commit(self, transaction).await? {
816                    guards.push(guard);
817                }
818            }
819        }
820
821        self.maybe_start_flush_task();
822
823        self.hooks().on_before_commit();
824
825        let _guard = debug_assert_not_too_long!(self.commit_mutex.lock());
826        let journal_offset = if self.journal().image_builder_mode().is_some() {
827            let journal_checkpoint =
828                JournalCheckpoint { file_offset: 0, checksum: 0, version: LATEST_VERSION };
829            let maybe_mutation = self
830                .object_manager()
831                .apply_transaction(transaction, &journal_checkpoint)
832                .expect("Transactions must not fail in image_builder_mode");
833            if let Some(mutation) = maybe_mutation {
834                assert!(matches!(mutation, Mutation::UpdateBorrowed(_)));
835                // These are Mutation::UpdateBorrowed which are normally used to track borrowing of
836                // metadata reservations. As we are image-building and not using the journal,
837                // we don't track this.
838            }
839            self.object_manager().did_commit_transaction(transaction, &journal_checkpoint, 0);
840            0
841        } else {
842            self.journal.commit(transaction).await?
843        };
844
845        std::mem::drop(guards);
846        self.completed_transactions.add(1);
847
848        // For now, call the callback whilst holding the lock.  Technically, we don't need to do
849        // that except if there's a post-commit-hook (which there usually won't be).  We can
850        // consider changing this if we need to for performance, but we'd need to double check that
851        // callers don't depend on this.
852        let result = callback(journal_offset);
853
854        if let Some(hook) = self.options.post_commit_hook.as_ref() {
855            hook().await;
856        }
857
858        Ok(result)
859    }
860
861    pub fn lock_manager(&self) -> &LockManager {
862        &self.lock_manager
863    }
864
865    pub fn hooks(&self) -> &Arc<HooksHandle> {
866        &self.options.hooks
867    }
868
869    pub(crate) fn drop_transaction(&self, transaction: &mut Transaction<'_>) {
870        self.objects.drop_transaction(transaction);
871        self.lock_manager.drop_transaction(transaction);
872    }
873
874    fn maybe_start_flush_task(&self) {
875        if self.journal.image_builder_mode().is_some() {
876            return;
877        }
878        let mut flush_task = self.flush_task.lock();
879        if flush_task.is_none() {
880            let journal = self.journal.clone();
881            *flush_task = Some(fasync::Task::spawn(
882                journal.flush_task().trace(trace_future_args!("Journal::flush_task")),
883            ));
884        }
885    }
886
887    fn start_trim_task(self: &Arc<Self>) {
888        if !self.device.supports_trim() {
889            info!("Device does not support trim; not scheduling trimming");
890            return;
891        }
892        let this = self.clone();
893        self.background_tasks
894            .spawn(this.trim_task().trace(trace_future_args!("Filesystem::trim_task")));
895    }
896
897    async fn trim_task(self: Arc<Self>) {
898        // This task will be cancelled when the filesystem is closed.
899        let Some((mut next_timer, _)) = self.options.trim_config else { return };
900        loop {
901            fasync::Timer::new(next_timer.clone()).await;
902
903            // The timer has fired indicating a trim is now due.  If we have a power manager, we
904            // now check to see if there's an external power source.
905            let start = Instant::now();
906            let result = if let Some(pm) = &self.options.power_manager {
907                let mut watcher = pm.clone().watch_battery();
908
909                // The pauser starts paused.
910                let pauser = Pauser::new(self.options.trim_charger_wait);
911
912                let mut pause_future = pin!(
913                    async {
914                        let mut wake_lease: Option<WakeLease> = None;
915                        loop {
916                            let Some((using_battery, new_lease)) = watcher.next_latest().await
917                            else {
918                                // If we lose the connection to the watcher, unpause and do not
919                                // worry about monitoring the power source.
920                                pauser.set_pause(false);
921                                drop(wake_lease); // Silence the compiler warnings.
922                                return;
923                            };
924
925                            // Pause if the device is using battery.
926                            pauser.set_pause(using_battery);
927
928                            // Hold onto a wake lease if we are using an external power source (and
929                            // we are therefore unpaused).
930                            if using_battery {
931                                wake_lease = None;
932                            } else if new_lease.is_some() {
933                                wake_lease = new_lease;
934                            }
935                        }
936                    }
937                    .fuse()
938                );
939
940                let mut do_trim = pin!(self.do_trim(Some(&pauser)).fuse());
941
942                loop {
943                    futures::select! {
944                        _ = pause_future => {}
945                        result = do_trim => break result,
946                    }
947                }
948
949                // Now that trim has completed, we don't need to watch the power source any more, so
950                // we just drop the pauser and the future monitoring the power source.
951            } else {
952                self.do_trim(None).await
953            };
954
955            let duration = start.elapsed();
956            match result {
957                Ok(bytes_trimmed) => info!(
958                    "Trimmed {bytes_trimmed} bytes in {duration:?}.  Next trim in \
959                     {next_timer:?}",
960                ),
961                Err(error) => error!(error:?; "Failed to trim"),
962            }
963
964            let Some((_, interval)) = self.options.trim_config else { return };
965            next_timer = interval;
966            if next_timer.is_zero() {
967                fasync::yield_now().await;
968            }
969        }
970    }
971
972    // Returns the number of bytes trimmed.
973    async fn do_trim(&self, pauser: Option<&Pauser>) -> Result<usize, Error> {
974        const MAX_EXTENTS_PER_BATCH: usize = 8;
975        const MAX_EXTENT_SIZE: usize = 256 * 1024;
976        let mut offset = 0;
977        let mut bytes_trimmed = 0;
978        loop {
979            let allocator = self.allocator();
980            if let Some(pauser) = pauser {
981                pauser.maybe_pause().await;
982            }
983            let trimmable_extents =
984                allocator.take_for_trimming(offset, MAX_EXTENT_SIZE, MAX_EXTENTS_PER_BATCH).await?;
985            for device_range in trimmable_extents.extents() {
986                self.device.trim(device_range.clone()).await?;
987                bytes_trimmed += device_range.length()? as usize;
988            }
989            if let Some(device_range) = trimmable_extents.extents().last() {
990                offset = device_range.end;
991            } else {
992                break;
993            }
994        }
995        Ok(bytes_trimmed)
996    }
997
998    fn start_clean_transfer_buffer_task(self: &Arc<Self>) {
999        let this = self.clone();
1000        self.background_tasks.spawn(
1001            async move {
1002                loop {
1003                    fasync::Timer::new(CLEAN_TRANSFER_BUFFER_INTERVAL).await;
1004                    this.device().clean_transfer_buffer();
1005                }
1006            }
1007            .trace(trace_future_args!("Filesystem::clean_transfer_buffer_task")),
1008        );
1009    }
1010
1011    /// Tops up the transaction's metadata_reservation to ensure the reservation is large enough to
1012    /// accommodate a maximally sized transaction.  This must be called before the transaction is
1013    /// used.
1014    ///
1015    /// If `skip_journal_checks` is unset, this function also ensures that there is sufficient
1016    /// space in the journal to write the transaction, and will block if the journal needs to grow.
1017    pub(crate) async fn add_transaction_reservation(
1018        self: &Arc<Self>,
1019        transaction: &mut Transaction<'_>,
1020    ) -> Result<(), Error> {
1021        if self.options.image_builder_mode.is_some() {
1022            // Image builder mode avoids the journal so reservation tracking for metadata overheads
1023            // doesn't make sense and so we essentially have 'all or nothing' semantics instead.
1024            return Ok(());
1025        }
1026        if !transaction.skip_journal_checks {
1027            self.maybe_start_flush_task();
1028            self.journal.check_journal_space().await?;
1029        }
1030
1031        match &mut transaction.metadata_reservation {
1032            MetadataReservation::BorrowedMetadata
1033            | MetadataReservation::BorrowedMetadataAndData => {}
1034            MetadataReservation::Hold(hold) => {
1035                let amount = TRANSACTION_METADATA_MAX_AMOUNT.saturating_sub(hold.amount());
1036                hold.add(hold.owner().reserve(amount).ok_or(FxfsError::NoSpace)?.forget());
1037            }
1038            MetadataReservation::Reservation(txn_reservation) => {
1039                txn_reservation.add(
1040                    self.allocator()
1041                        .reserve(
1042                            txn_reservation.owner_object_id(),
1043                            TRANSACTION_METADATA_MAX_AMOUNT
1044                                .saturating_sub(txn_reservation.amount()),
1045                        )
1046                        .ok_or(FxfsError::NoSpace)?
1047                        .forget(),
1048                );
1049            }
1050        }
1051        Ok(())
1052    }
1053
1054    pub(crate) async fn add_transaction(&self, skip_journal_checks: bool) {
1055        if skip_journal_checks {
1056            self.in_flight_transactions.fetch_add(1, Ordering::Relaxed);
1057        } else {
1058            let inc = || {
1059                let mut in_flights = self.in_flight_transactions.load(Ordering::Relaxed);
1060                while in_flights < MAX_IN_FLIGHT_TRANSACTIONS {
1061                    match self.in_flight_transactions.compare_exchange_weak(
1062                        in_flights,
1063                        in_flights + 1,
1064                        Ordering::Relaxed,
1065                        Ordering::Relaxed,
1066                    ) {
1067                        Ok(_) => return true,
1068                        Err(x) => in_flights = x,
1069                    }
1070                }
1071                return false;
1072            };
1073            while !inc() {
1074                let listener = self.transaction_limit_event.listen();
1075                if inc() {
1076                    break;
1077                }
1078                debug_assert_not_too_long!(listener);
1079            }
1080        }
1081    }
1082
1083    pub(crate) fn sub_transaction(&self) {
1084        let old = self.in_flight_transactions.fetch_sub(1, Ordering::Relaxed);
1085        assert!(old != 0);
1086        if old <= MAX_IN_FLIGHT_TRANSACTIONS {
1087            self.transaction_limit_event.notify(usize::MAX);
1088        }
1089    }
1090
1091    pub async fn truncate_guard(&self, store_id: u64, object_id: u64) -> TruncateGuard<'_> {
1092        let keys = lock_keys![LockKey::truncate(store_id, object_id,)];
1093        TruncateGuard(self.lock_manager().write_lock(keys).await)
1094    }
1095
1096    pub async fn truncate_guard_owned(
1097        self: &Arc<Self>,
1098        store_id: u64,
1099        object_id: u64,
1100    ) -> TruncateGuard<'static> {
1101        self.truncate_guard(store_id, object_id).await.into_owned(self.clone())
1102    }
1103
1104    /// Immediately tombstones (purges) an object in the graveyard, deallocating its extents,
1105    /// removing its graveyard entry, and inserting an LSM-tree tombstone record
1106    /// (`ObjectValue::None`) that deletes all remaining records for the object during compaction.
1107    ///
1108    /// Use this when space should be reclaimed immediately (e.g. when deleting a volume or
1109    /// unlinking a file with no open references) or when the background graveyard reaper is not
1110    /// running. Use [`Graveyard::queue_tombstone_object`] instead when tombstoning from a
1111    /// non-async context (such as `Drop`) or when work can be deferred to the background (such as
1112    /// initial mount-time reaping).
1113    pub async fn tombstone_object(
1114        &self,
1115        store_id: u64,
1116        object_id: u64,
1117        truncate_guard: Option<&TruncateGuard<'_>>,
1118    ) -> Result<(), Error> {
1119        let store = self
1120            .objects
1121            .store(store_id)
1122            .with_context(|| format!("Failed to get store {}", store_id))?;
1123        // For now, it's safe to assume that all objects in the root parent and root store should
1124        // return space to the metadata reservation, but we might have to revisit that if we end up
1125        // with objects that are in other stores.
1126        let options = if store_id == self.objects.root_parent_store_object_id()
1127            || store_id == self.objects.root_store_object_id()
1128        {
1129            transaction::Options {
1130                reservation: transaction::ReservationOptions::BorrowedMetadataAndData,
1131                ..Default::default()
1132            }
1133        } else {
1134            transaction::Options {
1135                reservation: transaction::ReservationOptions::BorrowedMetadata,
1136                ..Default::default()
1137            }
1138        };
1139        store.tombstone_object(object_id, options, truncate_guard).await
1140    }
1141
1142    /// Immediately tombstones (purges) an attribute in the graveyard, deallocating its extents,
1143    /// removing its graveyard entry, and inserting an LSM-tree tombstone record
1144    /// (`ObjectValue::None`) that deletes the attribute record during compaction.
1145    ///
1146    /// See [`FxFilesystem::tombstone_object`] for when to use this vs.
1147    /// [`Graveyard::queue_tombstone_attribute`].
1148    pub async fn tombstone_attribute(
1149        &self,
1150        store_id: u64,
1151        object_id: u64,
1152        attribute_id: AttributeId,
1153    ) -> Result<(), Error> {
1154        let store = self
1155            .objects
1156            .store(store_id)
1157            .with_context(|| format!("Failed to get store {}", store_id))?;
1158        // For now, it's safe to assume that all objects in the root parent and root store should
1159        // return space to the metadata reservation, but we might have to revisit that if we end up
1160        // with objects that are in other stores.
1161        let options = if store_id == self.objects.root_parent_store_object_id()
1162            || store_id == self.objects.root_store_object_id()
1163        {
1164            transaction::Options {
1165                reservation: transaction::ReservationOptions::BorrowedMetadataAndData,
1166                ..Default::default()
1167            }
1168        } else {
1169            transaction::Options {
1170                reservation: transaction::ReservationOptions::BorrowedMetadata,
1171                ..Default::default()
1172            }
1173        };
1174        store.tombstone_attribute(object_id, attribute_id, options).await
1175    }
1176
1177    async fn populate_stores_node(&self) -> Result<Inspector, Error> {
1178        let inspector = fuchsia_inspect::Inspector::default();
1179        let root = inspector.root();
1180        root.record_child("__root", |n| self.root_store().record_data(n));
1181        root.record_child("__root_parent", |n| self.root_parent_store().record_data(n));
1182        let object_manager = self.object_manager();
1183        let volume_directory = object_manager.volume_directory();
1184        let layer_set = volume_directory.store().tree().layer_set();
1185        let mut merger = layer_set.merger();
1186        let mut iter = volume_directory.iter(&mut merger).await?;
1187        while let Some((name, id, _)) = iter.get() {
1188            if let Some(store) = object_manager.store(id) {
1189                root.record_child(name.to_string(), |n| store.record_data(n));
1190            }
1191            iter.advance().await?;
1192        }
1193        Ok(inspector)
1194    }
1195}
1196
1197/// A wrapper around a guard that needs to be taken when truncating an object.
1198#[allow(dead_code)]
1199pub struct TruncateGuard<'a>(WriteGuard<'a>);
1200
1201impl<'a> TruncateGuard<'a> {
1202    pub fn into_owned(self, fs: Arc<FxFilesystem>) -> TruncateGuard<'static> {
1203        TruncateGuard(self.0.into_owned(fs))
1204    }
1205}
1206
1207/// Helper method for making a new filesystem.
1208pub async fn mkfs(device: DeviceHolder) -> Result<DeviceHolder, Error> {
1209    let fs = FxFilesystem::new_empty(device).await?;
1210    fs.close().await?;
1211    Ok(fs.take_device().await)
1212}
1213
1214/// Helper method for making a new filesystem with a single named volume.
1215/// This shouldn't be used in production; instead volumes should be created with the Volumes
1216/// protocol.
1217pub async fn mkfs_with_volume(
1218    device: DeviceHolder,
1219    volume_name: &str,
1220    crypt: Option<Arc<dyn Crypt>>,
1221) -> Result<DeviceHolder, Error> {
1222    let fs = FxFilesystem::new_empty(device).await?;
1223    {
1224        // expect instead of propagating errors here, since otherwise we could drop |fs| before
1225        // close is called, which leads to confusing and unrelated error messages.
1226        let root_volume = root_volume(fs.clone()).await.expect("Open root_volume failed");
1227        root_volume
1228            .new_volume(
1229                volume_name,
1230                NewChildStoreOptions {
1231                    options: StoreOptions { crypt, ..StoreOptions::default() },
1232                    ..Default::default()
1233                },
1234            )
1235            .await
1236            .expect("Create volume failed");
1237    }
1238    fs.close().await?;
1239    Ok(fs.take_device().await)
1240}
1241
1242struct FsckAfterEveryTransaction {
1243    fs: OnceLock<Weak<FxFilesystem>>,
1244    old_hook: PostCommitHook,
1245}
1246
1247impl FsckAfterEveryTransaction {
1248    fn new(old_hook: PostCommitHook) -> Arc<Self> {
1249        Arc::new(Self { fs: OnceLock::new(), old_hook })
1250    }
1251
1252    async fn run(self: Arc<Self>) {
1253        if let Some(fs) = self.fs.get().and_then(Weak::upgrade) {
1254            let options = FsckOptions {
1255                fail_on_warning: true,
1256                no_lock: true,
1257                quiet: true,
1258                ..Default::default()
1259            };
1260            fsck_with_options(fs.clone(), &options).await.expect("fsck failed");
1261            let object_manager = fs.object_manager();
1262            for store in object_manager.unlocked_stores() {
1263                let store_id = store.store_object_id();
1264                if !object_manager.is_system_store(store_id) {
1265                    fsck_volume_with_options(fs.as_ref(), &options, store_id, None)
1266                        .await
1267                        .expect("fsck_volume_with_options failed");
1268                }
1269            }
1270        }
1271        if let Some(old_hook) = self.old_hook.as_ref() {
1272            old_hook().await;
1273        }
1274    }
1275}
1276
1277struct Pauser {
1278    pause: Condition<bool>,
1279    bounce_delay: Duration,
1280}
1281
1282impl Pauser {
1283    /// Returns a new Pauser which starts paused.
1284    fn new(bounce_delay: Duration) -> Self {
1285        Self { pause: Condition::new(true), bounce_delay }
1286    }
1287
1288    async fn maybe_pause(&self) {
1289        loop {
1290            if !*self.pause.lock() {
1291                return;
1292            }
1293            self.pause.when(|p| if **p { Poll::Pending } else { Poll::Ready(()) }).await;
1294            fasync::Timer::new(self.bounce_delay).await;
1295        }
1296    }
1297
1298    fn set_pause(&self, v: bool) {
1299        let mut guard = self.pause.lock();
1300        if *guard == v {
1301            return;
1302        }
1303        *guard = v;
1304        for waker in guard.drain_wakers() {
1305            waker.wake();
1306        }
1307    }
1308}
1309
1310trait NextLatest: Stream + Unpin {
1311    /// Gets the next item from the stream, but if multiple items are ready, returns the latest.
1312    async fn next_latest(&mut self) -> Option<Self::Item> {
1313        let Some(mut next) = self.next().await else { return None };
1314
1315        // Coalesce with any subsequent items that are ready.
1316        loop {
1317            match self.next().now_or_never() {
1318                None => return Some(next),
1319                Some(None) => return None,
1320                Some(Some(n)) => next = n,
1321            }
1322        }
1323    }
1324}
1325
1326impl<T: ?Sized + Unpin> NextLatest for T where T: Stream {}
1327
1328#[cfg(test)]
1329mod tests {
1330    use super::{FxFilesystem, FxFilesystemBuilder, FxfsError, SyncOptions};
1331    use crate::fsck::{fsck, fsck_volume};
1332    use crate::log::*;
1333    use crate::lsm_tree::Operation;
1334    use crate::lsm_tree::types::Item;
1335    use crate::object_handle::{INVALID_OBJECT_ID, ObjectHandle, WriteObjectHandle};
1336    use crate::object_store::directory::{Directory, replace_child};
1337    use crate::object_store::journal::JournalOptions;
1338    use crate::object_store::journal::super_block::SuperBlockInstance;
1339    use crate::object_store::transaction::{LockKey, Options, lock_keys};
1340    use crate::object_store::volume::root_volume;
1341    use crate::object_store::{
1342        HandleOptions, NewChildStoreOptions, ObjectDescriptor, ObjectStore, StoreOptions,
1343    };
1344    use crate::range::RangeExt;
1345    use fuchsia_async as fasync;
1346    use fuchsia_sync::Mutex;
1347    use futures::future::join_all;
1348    use futures::stream::{FuturesUnordered, TryStreamExt};
1349    use fxfs_insecure_crypto::new_insecure_crypt;
1350    use rustc_hash::FxHashMap as HashMap;
1351    use std::ops::Range;
1352    use std::sync::Arc;
1353    use std::sync::atomic::{AtomicU32, Ordering};
1354    use std::time::Duration;
1355    use storage_device::DeviceHolder;
1356    use storage_device::fake_device::{self, FakeDevice};
1357    use test_case::test_case;
1358
1359    const TEST_DEVICE_BLOCK_SIZE: u32 = 512;
1360
1361    #[fuchsia::test(threads = 10)]
1362    async fn test_compaction() {
1363        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1364
1365        // If compaction is not working correctly, this test will run out of space.
1366        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1367        let root_store = fs.root_store();
1368        let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
1369            .await
1370            .expect("open failed");
1371
1372        let mut tasks = Vec::new();
1373        for i in 0..2 {
1374            let mut transaction = fs
1375                .root_store()
1376                .new_transaction(
1377                    lock_keys![LockKey::object(
1378                        root_store.store_object_id(),
1379                        root_directory.object_id()
1380                    )],
1381                    Options::default(),
1382                )
1383                .await
1384                .expect("new_transaction failed");
1385            let handle = root_directory
1386                .create_child_file(&mut transaction, &format!("{}", i))
1387                .await
1388                .expect("create_child_file failed");
1389            transaction.commit().await.expect("commit failed");
1390            tasks.push(fasync::Task::spawn(async move {
1391                const TEST_DATA: &[u8] = b"hello";
1392                let mut buf = handle.allocate_buffer(TEST_DATA.len()).await;
1393                buf.copy_from_slice(TEST_DATA);
1394                for _ in 0..1500 {
1395                    handle.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
1396                }
1397            }));
1398        }
1399        join_all(tasks).await;
1400        fs.sync(SyncOptions::default()).await.expect("sync failed");
1401
1402        fsck(fs.clone()).await.expect("fsck failed");
1403        fs.close().await.expect("Close failed");
1404    }
1405
1406    #[fuchsia::test]
1407    async fn test_enable_allocations() {
1408        // 1. enable_allocations() has no impact if image_builder_mode is not used.
1409        {
1410            let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1411            let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1412            fs.enable_allocations();
1413            let root_store = fs.root_store();
1414            let root_directory =
1415                Directory::open(&root_store, root_store.root_directory_object_id())
1416                    .await
1417                    .expect("open failed");
1418            let mut transaction = fs
1419                .root_store()
1420                .new_transaction(
1421                    lock_keys![LockKey::object(
1422                        root_store.store_object_id(),
1423                        root_directory.object_id()
1424                    )],
1425                    Options::default(),
1426                )
1427                .await
1428                .expect("new_transaction failed");
1429            root_directory
1430                .create_child_file(&mut transaction, "test")
1431                .await
1432                .expect("create_child_file failed");
1433            transaction.commit().await.expect("commit failed");
1434            fs.close().await.expect("close failed");
1435        }
1436
1437        // 2. Allocations blow up if done before this call (in image_builder_mode), but work after
1438        {
1439            let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1440            let fs = FxFilesystemBuilder::new()
1441                .format(true)
1442                .image_builder_mode(Some(SuperBlockInstance::A))
1443                .open(device)
1444                .await
1445                .expect("open failed");
1446            let root_store = fs.root_store();
1447            let root_directory =
1448                Directory::open(&root_store, root_store.root_directory_object_id())
1449                    .await
1450                    .expect("open failed");
1451
1452            let mut transaction = fs
1453                .root_store()
1454                .new_transaction(
1455                    lock_keys![LockKey::object(
1456                        root_store.store_object_id(),
1457                        root_directory.object_id()
1458                    )],
1459                    Options::default(),
1460                )
1461                .await
1462                .expect("new_transaction failed");
1463            let handle = root_directory
1464                .create_child_file(&mut transaction, "test_fail")
1465                .await
1466                .expect("create_child_file failed");
1467            transaction.commit().await.expect("commit failed");
1468
1469            // Allocations should fail before enable_allocations()
1470            assert!(
1471                FxfsError::Unavailable
1472                    .matches(&handle.allocate(0..4096).await.expect_err("allocate should fail"))
1473            );
1474
1475            // Allocations should work after enable_allocations()
1476            fs.enable_allocations();
1477            handle.allocate(0..4096).await.expect("allocate should work after enable_allocations");
1478
1479            // 3. finalize() works regardless of whether enable_allocations() is called.
1480            // (We already called it above, so this verifies it works after it was called).
1481
1482            fs.close().await.expect("close failed");
1483        }
1484        // TODO(https://fxbug.dev/467401079): Add a failure test where we close without
1485        // enabling allocations. (Trivial to do, but causes error logs, which are interpreted as
1486        // test failures and only seem controllable at the BUILD target level).
1487    }
1488
1489    #[fuchsia::test(threads = 10)]
1490    async fn test_replay_is_identical() {
1491        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1492        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1493
1494        // Reopen the store, but set reclaim size to a very large value which will effectively
1495        // stop the journal from flushing and allows us to track all the mutations to the store.
1496        fs.close().await.expect("close failed");
1497        let device = fs.take_device().await;
1498        device.reopen(false);
1499
1500        struct Mutations<K, V>(Mutex<Vec<(Operation, Item<K, V>)>>);
1501
1502        impl<K: Clone, V: Clone> Mutations<K, V> {
1503            fn new() -> Self {
1504                Mutations(Mutex::new(Vec::new()))
1505            }
1506
1507            fn push(&self, operation: Operation, item: &Item<K, V>) {
1508                self.0.lock().push((operation, item.clone()));
1509            }
1510        }
1511
1512        let open_fs = |device,
1513                       object_mutations: Arc<Mutex<HashMap<_, _>>>,
1514                       allocator_mutations: Arc<Mutations<_, _>>| async {
1515            FxFilesystemBuilder::new()
1516                .journal_options(JournalOptions { reclaim_size: u64::MAX, ..Default::default() })
1517                .on_new_allocator(move |allocator| {
1518                    let allocator_mutations = allocator_mutations.clone();
1519                    allocator.tree().set_mutation_callback(Some(Box::new(move |op, item| {
1520                        allocator_mutations.push(op, item)
1521                    })));
1522                })
1523                .on_new_store(move |store| {
1524                    let mutations = Arc::new(Mutations::new());
1525                    object_mutations.lock().insert(store.store_object_id(), mutations.clone());
1526                    store.tree().set_mutation_callback(Some(Box::new(move |op, item| {
1527                        mutations.push(op, item)
1528                    })));
1529                })
1530                .open(device)
1531                .await
1532                .expect("open failed")
1533        };
1534
1535        let allocator_mutations = Arc::new(Mutations::new());
1536        let object_mutations = Arc::new(Mutex::new(HashMap::default()));
1537        let fs = open_fs(device, object_mutations.clone(), allocator_mutations.clone()).await;
1538
1539        let root_store = fs.root_store();
1540        let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
1541            .await
1542            .expect("open failed");
1543
1544        let mut transaction = fs
1545            .root_store()
1546            .new_transaction(
1547                lock_keys![LockKey::object(
1548                    root_store.store_object_id(),
1549                    root_directory.object_id()
1550                )],
1551                Options::default(),
1552            )
1553            .await
1554            .expect("new_transaction failed");
1555        let object = root_directory
1556            .create_child_file(&mut transaction, "test")
1557            .await
1558            .expect("create_child_file failed");
1559        transaction.commit().await.expect("commit failed");
1560
1561        // Append some data.
1562        let buf = object.allocate_buffer(10000).await;
1563        object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
1564
1565        // Overwrite some data.
1566        object.write_or_append(Some(5000), buf.as_ref()).await.expect("write failed");
1567
1568        // Truncate.
1569        object.truncate(3000).await.expect("truncate failed");
1570
1571        // Delete the object.
1572        let mut transaction = fs
1573            .root_store()
1574            .new_transaction(
1575                lock_keys![
1576                    LockKey::object(root_store.store_object_id(), root_directory.object_id()),
1577                    LockKey::object(root_store.store_object_id(), object.object_id()),
1578                ],
1579                Options::default(),
1580            )
1581            .await
1582            .expect("new_transaction failed");
1583
1584        replace_child(&mut transaction, None, (&root_directory, "test"))
1585            .await
1586            .expect("replace_child failed");
1587
1588        transaction.commit().await.expect("commit failed");
1589
1590        // Finally tombstone the object.
1591        root_store
1592            .tombstone_object(object.object_id(), Options::default(), None)
1593            .await
1594            .expect("tombstone failed");
1595
1596        // Now reopen and check that replay produces the same set of mutations.
1597        fs.close().await.expect("close failed");
1598
1599        let metadata_reservation_amount = fs.object_manager().metadata_reservation().amount();
1600
1601        let device = fs.take_device().await;
1602        device.reopen(false);
1603
1604        let replayed_object_mutations = Arc::new(Mutex::new(HashMap::default()));
1605        let replayed_allocator_mutations = Arc::new(Mutations::new());
1606        let fs = open_fs(
1607            device,
1608            replayed_object_mutations.clone(),
1609            replayed_allocator_mutations.clone(),
1610        )
1611        .await;
1612
1613        let m1 = object_mutations.lock();
1614        let m2 = replayed_object_mutations.lock();
1615        assert_eq!(m1.len(), m2.len());
1616        for (store_id, mutations) in &*m1 {
1617            let mutations = mutations.0.lock();
1618            let replayed = m2.get(&store_id).expect("Found unexpected store").0.lock();
1619            assert_eq!(mutations.len(), replayed.len());
1620            for ((op1, i1), (op2, i2)) in mutations.iter().zip(replayed.iter()) {
1621                assert_eq!(op1, op2);
1622                assert_eq!(i1.key, i2.key);
1623                assert_eq!(i1.value, i2.value);
1624            }
1625        }
1626
1627        let a1 = allocator_mutations.0.lock();
1628        let a2 = replayed_allocator_mutations.0.lock();
1629        assert_eq!(a1.len(), a2.len());
1630        for ((op1, i1), (op2, i2)) in a1.iter().zip(a2.iter()) {
1631            assert_eq!(op1, op2);
1632            assert_eq!(i1.key, i2.key);
1633            assert_eq!(i1.value, i2.value);
1634        }
1635
1636        assert_eq!(
1637            fs.object_manager().metadata_reservation().amount(),
1638            metadata_reservation_amount
1639        );
1640    }
1641
1642    #[fuchsia::test]
1643    async fn test_max_in_flight_transactions() {
1644        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1645        let fs = FxFilesystem::new_empty(device).await.expect("new_empty failed");
1646
1647        let store = fs.root_store();
1648        let transactions = FuturesUnordered::new();
1649        for _ in 0..super::MAX_IN_FLIGHT_TRANSACTIONS {
1650            transactions.push(store.new_transaction(lock_keys![], Options::default()));
1651        }
1652        let mut transactions: Vec<_> = transactions.try_collect().await.unwrap();
1653
1654        // Trying to create another one should be blocked.
1655        let mut fut = std::pin::pin!(store.new_transaction(lock_keys![], Options::default()));
1656        assert!(futures::poll!(&mut fut).is_pending());
1657
1658        // Dropping one should allow it to proceed.
1659        transactions.pop();
1660
1661        assert!(futures::poll!(&mut fut).is_ready());
1662    }
1663
1664    // If run on a single thread, the trim tasks starve out other work.
1665    #[fuchsia::test(threads = 10)]
1666    async fn test_continuously_trim() {
1667        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
1668        let fs = FxFilesystemBuilder::new()
1669            .trim_config(Some((Duration::ZERO, Duration::ZERO)))
1670            .format(true)
1671            .open(device)
1672            .await
1673            .expect("open failed");
1674        // Do a small sleep so trim has time to get going.
1675        fasync::Timer::new(Duration::from_millis(10)).await;
1676
1677        // Create and delete a bunch of files whilst trim is ongoing.  This just ensures that
1678        // regular usage isn't affected by trim.
1679        let root_store = fs.root_store();
1680        let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
1681            .await
1682            .expect("open failed");
1683        for _ in 0..100 {
1684            let mut transaction = fs
1685                .root_store()
1686                .new_transaction(
1687                    lock_keys![LockKey::object(
1688                        root_store.store_object_id(),
1689                        root_directory.object_id()
1690                    )],
1691                    Options::default(),
1692                )
1693                .await
1694                .expect("new_transaction failed");
1695            let object = root_directory
1696                .create_child_file(&mut transaction, "test")
1697                .await
1698                .expect("create_child_file failed");
1699            transaction.commit().await.expect("commit failed");
1700
1701            {
1702                let buf = object.allocate_buffer(1024).await;
1703                object.write_or_append(Some(0), buf.as_ref()).await.expect("write failed");
1704            }
1705            std::mem::drop(object);
1706
1707            let mut transaction = root_directory
1708                .acquire_context_for_replace(None, "test", true)
1709                .await
1710                .expect("acquire_context_for_replace failed")
1711                .transaction;
1712            replace_child(&mut transaction, None, (&root_directory, "test"))
1713                .await
1714                .expect("replace_child failed");
1715            transaction.commit().await.expect("commit failed");
1716        }
1717        fs.close().await.expect("close failed");
1718    }
1719
1720    #[test_case(true; "test power fail with barriers")]
1721    #[test_case(false; "test power fail with checksums")]
1722    #[fuchsia::test]
1723    async fn test_power_fail(barriers_enabled: bool) {
1724        // This test randomly discards blocks, so we run it a few times to increase the chances
1725        // of catching an issue in a single run.
1726        for _ in 0..10 {
1727            let (store_id, device, test_file_object_id) = {
1728                let device = DeviceHolder::new(FakeDevice::new(8192, 4096));
1729                let fs = if barriers_enabled {
1730                    FxFilesystemBuilder::new()
1731                        .barriers_enabled(true)
1732                        .format(true)
1733                        .open(device)
1734                        .await
1735                        .expect("new filesystem failed")
1736                } else {
1737                    FxFilesystem::new_empty(device).await.expect("new_empty failed")
1738                };
1739                let root_volume = root_volume(fs.clone()).await.expect("root_volume failed");
1740
1741                fs.sync(SyncOptions { flush_device: true, ..SyncOptions::default() })
1742                    .await
1743                    .expect("sync failed");
1744
1745                let store = root_volume
1746                    .new_volume(
1747                        "test",
1748                        NewChildStoreOptions {
1749                            options: StoreOptions {
1750                                crypt: Some(Arc::new(new_insecure_crypt())),
1751                                ..StoreOptions::default()
1752                            },
1753                            ..Default::default()
1754                        },
1755                    )
1756                    .await
1757                    .expect("new_volume failed");
1758                let root_directory = Directory::open(&store, store.root_directory_object_id())
1759                    .await
1760                    .expect("open failed");
1761
1762                // Create a number of files with the goal of using up more than one journal block.
1763                async fn create_files(store: &Arc<ObjectStore>, prefix: &str) {
1764                    let fs = store.filesystem();
1765                    let root_directory = Directory::open(store, store.root_directory_object_id())
1766                        .await
1767                        .expect("open failed");
1768                    for i in 0..100 {
1769                        let mut transaction = fs
1770                            .root_store()
1771                            .new_transaction(
1772                                lock_keys![LockKey::object(
1773                                    store.store_object_id(),
1774                                    store.root_directory_object_id()
1775                                )],
1776                                Options::default(),
1777                            )
1778                            .await
1779                            .expect("new_transaction failed");
1780                        root_directory
1781                            .create_child_file(&mut transaction, &format!("{prefix} {i}"))
1782                            .await
1783                            .expect("create_child_file failed");
1784                        transaction.commit().await.expect("commit failed");
1785                    }
1786                }
1787
1788                // Create one batch of files.
1789                create_files(&store, "A").await;
1790
1791                // Create a file and write something to it.  This will make sure there's a
1792                // transaction present that includes a checksum.
1793                let mut transaction = fs
1794                    .root_store()
1795                    .new_transaction(
1796                        lock_keys![LockKey::object(
1797                            store.store_object_id(),
1798                            store.root_directory_object_id()
1799                        )],
1800                        Options::default(),
1801                    )
1802                    .await
1803                    .expect("new_transaction failed");
1804                let object = root_directory
1805                    .create_child_file(&mut transaction, "test")
1806                    .await
1807                    .expect("create_child_file failed");
1808                transaction.commit().await.expect("commit failed");
1809
1810                let mut transaction =
1811                    object.new_transaction().await.expect("new_transaction failed");
1812                let mut buffer = object.allocate_buffer(4096).await;
1813                buffer.fill(0xed);
1814                object
1815                    .txn_write(&mut transaction, 0, buffer.as_ref())
1816                    .await
1817                    .expect("txn_write failed");
1818                transaction.commit().await.expect("commit failed");
1819
1820                // Create another batch of files.
1821                create_files(&store, "B").await;
1822
1823                // Sync the device, but don't flush the device. We want to do this so we can
1824                // randomly discard blocks below.
1825                fs.sync(SyncOptions::default()).await.expect("sync failed");
1826
1827                // When we call `sync` above on the filesystem, it will pad the journal so that it
1828                // will get written, but it doesn't wait for the write to occur.  We wait for a
1829                // short time here to give allow time for the journal to be written.  Adding timers
1830                // isn't great, but this test already isn't deterministic since we randomly discard
1831                // blocks.
1832                fasync::Timer::new(Duration::from_millis(10)).await;
1833
1834                (
1835                    store.store_object_id(),
1836                    fs.device().snapshot().expect("snapshot failed"),
1837                    object.object_id(),
1838                )
1839            };
1840
1841            // Randomly discard blocks since the last flush.  This simulates what might happen in
1842            // the case of power-loss.  This will be an uncontrolled unmount.
1843            device
1844                .discard_random_since_last_flush()
1845                .expect("discard_random_since_last_flush failed");
1846
1847            let fs = FxFilesystem::open(device).await.expect("open failed");
1848            fsck(fs.clone()).await.expect("fsck failed");
1849
1850            let mut check_test_file = false;
1851
1852            // If we replayed and the store exists (i.e. the transaction that created the store
1853            // made it out), start by running fsck on it.
1854            let object_id = if fs.object_manager().store(store_id).is_some() {
1855                fsck_volume(&fs, store_id, Some(Arc::new(new_insecure_crypt())))
1856                    .await
1857                    .expect("fsck_volume failed");
1858
1859                // Now we want to create another file, unmount cleanly, and then finally check that
1860                // the new file exists.  This checks that we can continue to use the filesystem
1861                // after an unclean unmount.
1862                let store = root_volume(fs.clone())
1863                    .await
1864                    .expect("root_volume failed")
1865                    .volume(
1866                        "test",
1867                        StoreOptions {
1868                            crypt: Some(Arc::new(new_insecure_crypt())),
1869                            ..StoreOptions::default()
1870                        },
1871                    )
1872                    .await
1873                    .expect("volume failed");
1874
1875                let root_directory = Directory::open(&store, store.root_directory_object_id())
1876                    .await
1877                    .expect("open failed");
1878
1879                let mut transaction = fs
1880                    .root_store()
1881                    .new_transaction(
1882                        lock_keys![LockKey::object(
1883                            store.store_object_id(),
1884                            store.root_directory_object_id()
1885                        )],
1886                        Options::default(),
1887                    )
1888                    .await
1889                    .expect("new_transaction failed");
1890                let object = root_directory
1891                    .create_child_file(&mut transaction, &format!("C"))
1892                    .await
1893                    .expect("create_child_file failed");
1894                transaction.commit().await.expect("commit failed");
1895
1896                // Write again to the test file if it exists.
1897                if let Ok(test_file) = ObjectStore::open_object(
1898                    &store,
1899                    test_file_object_id,
1900                    HandleOptions::default(),
1901                    None,
1902                )
1903                .await
1904                {
1905                    // Check it has the contents we expect.
1906                    let mut buffer = test_file.allocate_buffer(4096).await;
1907                    let bytes =
1908                        test_file.read_aligned(0, buffer.as_mut()).await.expect("read failed");
1909                    if bytes == 4096 {
1910                        let expected = [0xed; 4096];
1911                        assert_eq!(buffer.to_vec(), expected);
1912                    } else {
1913                        // If the write didn't make it, the file should have zero bytes.
1914                        assert_eq!(bytes, 0);
1915                    }
1916
1917                    // Modify the test file.
1918                    let mut transaction =
1919                        test_file.new_transaction().await.expect("new_transaction failed");
1920                    buffer.fill(0x37);
1921                    test_file
1922                        .txn_write(&mut transaction, 0, buffer.as_ref())
1923                        .await
1924                        .expect("txn_write failed");
1925                    transaction.commit().await.expect("commit failed");
1926                    check_test_file = true;
1927                }
1928
1929                object.object_id()
1930            } else {
1931                INVALID_OBJECT_ID
1932            };
1933
1934            // This will do a controlled unmount.
1935            fs.close().await.expect("close failed");
1936            let device = fs.take_device().await;
1937            device.reopen(false);
1938
1939            let fs = FxFilesystem::open(device).await.expect("open failed");
1940            fsck(fs.clone()).await.expect("fsck failed");
1941
1942            // As mentioned above, make sure that the object we created before the clean unmount
1943            // exists.
1944            if object_id != INVALID_OBJECT_ID {
1945                fsck_volume(&fs, store_id, Some(Arc::new(new_insecure_crypt())))
1946                    .await
1947                    .expect("fsck_volume failed");
1948
1949                let store = root_volume(fs.clone())
1950                    .await
1951                    .expect("root_volume failed")
1952                    .volume(
1953                        "test",
1954                        StoreOptions {
1955                            crypt: Some(Arc::new(new_insecure_crypt())),
1956                            ..StoreOptions::default()
1957                        },
1958                    )
1959                    .await
1960                    .expect("volume failed");
1961                // We should be able to open the C object.
1962                ObjectStore::open_object(&store, object_id, HandleOptions::default(), None)
1963                    .await
1964                    .expect("open_object failed");
1965
1966                // If we made the modification to the test file, check it.
1967                if check_test_file {
1968                    info!("Checking test file for modification");
1969                    let test_file = ObjectStore::open_object(
1970                        &store,
1971                        test_file_object_id,
1972                        HandleOptions::default(),
1973                        None,
1974                    )
1975                    .await
1976                    .expect("open_object failed");
1977                    let mut buffer = test_file.allocate_buffer(4096).await;
1978                    assert_eq!(
1979                        test_file.read_aligned(0, buffer.as_mut()).await.expect("read failed"),
1980                        4096
1981                    );
1982                    let expected = [0x37; 4096];
1983                    let data = buffer.to_vec();
1984                    assert_eq!(data, expected);
1985                }
1986            }
1987
1988            fs.close().await.expect("close failed");
1989        }
1990    }
1991
1992    #[fuchsia::test]
1993    async fn test_barrier_not_emitted_when_transaction_has_no_data() {
1994        let barrier_count = Arc::new(AtomicU32::new(0));
1995
1996        struct Observer(Arc<AtomicU32>);
1997
1998        impl fake_device::Observer for Observer {
1999            fn barrier(&self) {
2000                self.0.fetch_add(1, Ordering::Relaxed);
2001            }
2002        }
2003
2004        let mut fake_device = FakeDevice::new(8192, 4096);
2005        fake_device.set_observer(Box::new(Observer(barrier_count.clone())));
2006        let device = DeviceHolder::new(fake_device);
2007        let fs = FxFilesystemBuilder::new()
2008            .barriers_enabled(true)
2009            .format(true)
2010            .open(device)
2011            .await
2012            .expect("new filesystem failed");
2013
2014        {
2015            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2016            root_vol
2017                .new_volume(
2018                    "test",
2019                    NewChildStoreOptions {
2020                        options: StoreOptions {
2021                            crypt: Some(Arc::new(new_insecure_crypt())),
2022                            ..StoreOptions::default()
2023                        },
2024                        ..NewChildStoreOptions::default()
2025                    },
2026                )
2027                .await
2028                .expect("there is no test volume");
2029            fs.close().await.expect("close failed");
2030        }
2031        // Remount the filesystem to ensure that the journal flushes and we can get a reliable
2032        // measure of the number of barriers issued during setup.
2033        let device = fs.take_device().await;
2034        device.reopen(false);
2035        let fs = FxFilesystemBuilder::new()
2036            .barriers_enabled(true)
2037            .open(device)
2038            .await
2039            .expect("new filesystem failed");
2040        let expected_barrier_count = barrier_count.load(Ordering::Relaxed);
2041
2042        let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2043        let store = root_vol
2044            .volume(
2045                "test",
2046                StoreOptions {
2047                    crypt: Some(Arc::new(new_insecure_crypt())),
2048                    ..StoreOptions::default()
2049                },
2050            )
2051            .await
2052            .expect("there is no test volume");
2053
2054        // Create a number of files with the goal of using up more than one journal block.
2055        let fs = store.filesystem();
2056        let root_directory =
2057            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
2058        for i in 0..100 {
2059            let mut transaction = fs
2060                .root_store()
2061                .new_transaction(
2062                    lock_keys![LockKey::object(
2063                        store.store_object_id(),
2064                        store.root_directory_object_id()
2065                    )],
2066                    Options::default(),
2067                )
2068                .await
2069                .expect("new_transaction failed");
2070            root_directory
2071                .create_child_file(&mut transaction, &format!("A {i}"))
2072                .await
2073                .expect("create_child_file failed");
2074            transaction.commit().await.expect("commit failed");
2075        }
2076
2077        // Unmount the filesystem to ensure that the journal flushes.
2078        fs.close().await.expect("close failed");
2079        // Ensure that no barriers were emitted while creating files, as no data was written.
2080        assert_eq!(expected_barrier_count, barrier_count.load(Ordering::Relaxed));
2081    }
2082
2083    #[fuchsia::test]
2084    async fn test_barrier_emitted_when_transaction_includes_data() {
2085        let barrier_count = Arc::new(AtomicU32::new(0));
2086
2087        struct Observer(Arc<AtomicU32>);
2088
2089        impl fake_device::Observer for Observer {
2090            fn barrier(&self) {
2091                self.0.fetch_add(1, Ordering::Relaxed);
2092            }
2093        }
2094
2095        let mut fake_device = FakeDevice::new(8192, 4096);
2096        fake_device.set_observer(Box::new(Observer(barrier_count.clone())));
2097        let device = DeviceHolder::new(fake_device);
2098        let fs = FxFilesystemBuilder::new()
2099            .barriers_enabled(true)
2100            .format(true)
2101            .open(device)
2102            .await
2103            .expect("new filesystem failed");
2104
2105        {
2106            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2107            root_vol
2108                .new_volume(
2109                    "test",
2110                    NewChildStoreOptions {
2111                        options: StoreOptions {
2112                            crypt: Some(Arc::new(new_insecure_crypt())),
2113                            ..StoreOptions::default()
2114                        },
2115                        ..NewChildStoreOptions::default()
2116                    },
2117                )
2118                .await
2119                .expect("there is no test volume");
2120            fs.close().await.expect("close failed");
2121        }
2122        // Remount the filesystem to ensure that the journal flushes and we can get a reliable
2123        // measure of the number of barriers issued during setup.
2124        let device = fs.take_device().await;
2125        device.reopen(false);
2126        let fs = FxFilesystemBuilder::new()
2127            .barriers_enabled(true)
2128            .open(device)
2129            .await
2130            .expect("new filesystem failed");
2131        let expected_barrier_count = barrier_count.load(Ordering::Relaxed);
2132
2133        let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2134        let store = root_vol
2135            .volume(
2136                "test",
2137                StoreOptions {
2138                    crypt: Some(Arc::new(new_insecure_crypt())),
2139                    ..StoreOptions::default()
2140                },
2141            )
2142            .await
2143            .expect("there is no test volume");
2144
2145        // Create a file and write something to it. This should cause a barrier to be emitted.
2146        let fs: Arc<FxFilesystem> = store.filesystem();
2147        let root_directory =
2148            Directory::open(&store, store.root_directory_object_id()).await.expect("open failed");
2149
2150        let mut transaction = fs
2151            .root_store()
2152            .new_transaction(
2153                lock_keys![LockKey::object(
2154                    store.store_object_id(),
2155                    store.root_directory_object_id()
2156                )],
2157                Options::default(),
2158            )
2159            .await
2160            .expect("new_transaction failed");
2161        let object = root_directory
2162            .create_child_file(&mut transaction, "test")
2163            .await
2164            .expect("create_child_file failed");
2165        transaction.commit().await.expect("commit failed");
2166
2167        let mut transaction = object.new_transaction().await.expect("new_transaction failed");
2168        let mut buffer = object.allocate_buffer(4096).await;
2169        buffer.fill(0xed);
2170        object.txn_write(&mut transaction, 0, buffer.as_ref()).await.expect("txn_write failed");
2171        transaction.commit().await.expect("commit failed");
2172
2173        // Unmount the filesystem to ensure that the journal flushes.
2174        fs.close().await.expect("close failed");
2175        // Ensure that a barrier was emitted while writing to the file.
2176        assert!(expected_barrier_count < barrier_count.load(Ordering::Relaxed));
2177    }
2178
2179    #[test_case(true; "fail when original filesystem has barriers enabled")]
2180    #[test_case(false; "fail when original filesystem has barriers disabled")]
2181    #[fuchsia::test]
2182    async fn test_switching_barrier_mode_on_existing_filesystem(original_barrier_mode: bool) {
2183        let crypt = Some(Arc::new(new_insecure_crypt()) as Arc<dyn fxfs_crypto::Crypt>);
2184        let fake_device = FakeDevice::new(8192, 4096);
2185        let device = DeviceHolder::new(fake_device);
2186        let fs: super::OpenFxFilesystem = FxFilesystemBuilder::new()
2187            .barriers_enabled(original_barrier_mode)
2188            .format(true)
2189            .open(device)
2190            .await
2191            .expect("new filesystem failed");
2192
2193        // Create a volume named test with a file inside it called file.
2194        {
2195            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2196            let store = root_vol
2197                .new_volume(
2198                    "test",
2199                    NewChildStoreOptions {
2200                        options: StoreOptions { crypt: crypt.clone(), ..Default::default() },
2201                        ..Default::default()
2202                    },
2203                )
2204                .await
2205                .expect("creating test volume");
2206            let root_dir = Directory::open(&store, store.root_directory_object_id())
2207                .await
2208                .expect("open failed");
2209            let mut transaction = fs
2210                .root_store()
2211                .new_transaction(
2212                    lock_keys![LockKey::object(
2213                        store.store_object_id(),
2214                        store.root_directory_object_id()
2215                    )],
2216                    Default::default(),
2217                )
2218                .await
2219                .expect("new_transaction failed");
2220            let object = root_dir
2221                .create_child_file(&mut transaction, "file")
2222                .await
2223                .expect("create_child_file failed");
2224            transaction.commit().await.expect("commit failed");
2225            let mut buffer = object.allocate_buffer(4096).await;
2226            buffer.fill(0xA7);
2227            let new_size = object.write_or_append(None, buffer.as_ref()).await.unwrap();
2228            assert_eq!(new_size, 4096);
2229        }
2230
2231        // Remount the filesystem with the opposite barrier mode and write more data to our file.
2232        fs.close().await.expect("close failed");
2233        let device = fs.take_device().await;
2234        device.reopen(false);
2235        let fs = FxFilesystemBuilder::new()
2236            .barriers_enabled(!original_barrier_mode)
2237            .open(device)
2238            .await
2239            .expect("new filesystem failed");
2240        {
2241            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2242            let store = root_vol
2243                .volume("test", StoreOptions { crypt: crypt.clone(), ..Default::default() })
2244                .await
2245                .expect("opening test volume");
2246            let root_dir = Directory::open(&store, store.root_directory_object_id())
2247                .await
2248                .expect("open failed");
2249            let (object_id, _, _) =
2250                root_dir.lookup("file").await.expect("lookup failed").expect("missing file");
2251            let test_file = ObjectStore::open_object(&store, object_id, Default::default(), None)
2252                .await
2253                .expect("open failed");
2254            // Write some more data.
2255            let mut buffer = test_file.allocate_buffer(4096).await;
2256            buffer.fill(0xA8);
2257            let new_size = test_file.write_or_append(None, buffer.as_ref()).await.unwrap();
2258            assert_eq!(new_size, 8192);
2259        }
2260
2261        // Lastly, remount the filesystems with the original barrier mode and make sure everything
2262        // can be read from the file as expected.
2263        fs.close().await.expect("close failed");
2264        let device = fs.take_device().await;
2265        device.reopen(false);
2266        let fs = FxFilesystemBuilder::new()
2267            .barriers_enabled(original_barrier_mode)
2268            .open(device)
2269            .await
2270            .expect("new filesystem failed");
2271        {
2272            let root_vol = root_volume(fs.clone()).await.expect("root_volume failed");
2273            let store = root_vol
2274                .volume("test", StoreOptions { crypt: crypt.clone(), ..Default::default() })
2275                .await
2276                .expect("opening test volume");
2277            let root_dir = Directory::open(&store, store.root_directory_object_id())
2278                .await
2279                .expect("open failed");
2280            let (object_id, _, _) =
2281                root_dir.lookup("file").await.expect("lookup failed").expect("missing file");
2282            let test_file = ObjectStore::open_object(&store, object_id, Default::default(), None)
2283                .await
2284                .expect("open failed");
2285            let mut buffer = test_file.allocate_buffer(8192).await;
2286            assert_eq!(
2287                test_file.read_aligned(0, buffer.as_mut()).await.expect("read failed"),
2288                8192,
2289                "short read"
2290            );
2291            let data = buffer.to_vec();
2292            assert_eq!(data[0..4096], [0xA7; 4096]);
2293            assert_eq!(data[4096..8192], [0xA8; 4096]);
2294        }
2295        fs.close().await.expect("close failed");
2296    }
2297
2298    #[fuchsia::test]
2299    async fn test_image_builder_mode_no_early_writes() {
2300        const BLOCK_SIZE: u32 = 4096;
2301        let device = DeviceHolder::new(FakeDevice::new(2048, BLOCK_SIZE));
2302        device.reopen(true);
2303        let fs = FxFilesystemBuilder::new()
2304            .format(true)
2305            .image_builder_mode(Some(SuperBlockInstance::A))
2306            .open(device)
2307            .await
2308            .expect("open failed");
2309        fs.enable_allocations();
2310        // fs.close() now performs compaction (writing superblock), so device must be writable.
2311        fs.device().reopen(false);
2312        fs.close().await.expect("closed");
2313    }
2314
2315    #[fuchsia::test]
2316    async fn test_image_builder_mode() {
2317        const BLOCK_SIZE: u32 = 4096;
2318        const EXISTING_FILE_RANGE: Range<u64> = 4096 * 1024..4096 * 1025;
2319        let device = DeviceHolder::new(FakeDevice::new(2048, BLOCK_SIZE));
2320
2321        // Write some fake file data at an offset in the image and confirm it as an fxfs file below.
2322        {
2323            let mut write_buf =
2324                device.allocate_buffer(EXISTING_FILE_RANGE.length().unwrap() as usize).await;
2325            write_buf.fill(0xf0);
2326            device.write(EXISTING_FILE_RANGE.start, write_buf.as_ref()).await.expect("write");
2327        }
2328
2329        device.reopen(true);
2330
2331        let device = {
2332            let fs = FxFilesystemBuilder::new()
2333                .format(true)
2334                .image_builder_mode(Some(SuperBlockInstance::B))
2335                .open(device)
2336                .await
2337                .expect("open failed");
2338            fs.enable_allocations();
2339            {
2340                let root_store = fs.root_store();
2341                let root_directory =
2342                    Directory::open(&root_store, root_store.root_directory_object_id())
2343                        .await
2344                        .expect("open failed");
2345                // Create a file referencing existing data on device.
2346                let handle;
2347                {
2348                    let mut transaction = fs
2349                        .root_store()
2350                        .new_transaction(
2351                            lock_keys![LockKey::object(
2352                                root_directory.store().store_object_id(),
2353                                root_directory.object_id()
2354                            )],
2355                            Options::default(),
2356                        )
2357                        .await
2358                        .expect("new transaction");
2359                    handle = root_directory
2360                        .create_child_file(&mut transaction, "test")
2361                        .await
2362                        .expect("create file");
2363                    handle.extend(&mut transaction, EXISTING_FILE_RANGE).await.expect("extend");
2364                    transaction.commit().await.expect("commit");
2365                }
2366            }
2367            fs.device().reopen(false);
2368            fs.close().await.expect("close");
2369            fs.take_device().await
2370        };
2371        device.reopen(false);
2372        let fs = FxFilesystem::open(device).await.expect("open failed");
2373        fsck(fs.clone()).await.expect("fsck failed");
2374
2375        // Confirm that the test file points at the correct data.
2376        let root_store = fs.root_store();
2377        let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
2378            .await
2379            .expect("open failed");
2380        let (object_id, descriptor, _) =
2381            root_directory.lookup("test").await.expect("lookup failed").unwrap();
2382        assert_eq!(descriptor, ObjectDescriptor::File);
2383        let test_file =
2384            ObjectStore::open_object(&root_store, object_id, HandleOptions::default(), None)
2385                .await
2386                .expect("open failed");
2387        let mut read_buf =
2388            test_file.allocate_buffer(EXISTING_FILE_RANGE.length().unwrap() as usize).await;
2389        test_file.read_aligned(0, read_buf.as_mut()).await.expect("read failed");
2390        let data = read_buf.to_vec();
2391        assert_eq!(data, [0xf0; 4096]);
2392        fs.close().await.expect("closed");
2393    }
2394
2395    #[fuchsia::test]
2396    async fn test_read_only_mount_on_full_filesystem() {
2397        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
2398        let fs =
2399            FxFilesystemBuilder::new().format(true).open(device).await.expect("new_empty failed");
2400        let root_store = fs.root_store();
2401        let root_directory = Directory::open(&root_store, root_store.root_directory_object_id())
2402            .await
2403            .expect("open failed");
2404
2405        let mut transaction = fs
2406            .root_store()
2407            .new_transaction(
2408                lock_keys![LockKey::object(
2409                    root_store.store_object_id(),
2410                    root_directory.object_id()
2411                )],
2412                Options::default(),
2413            )
2414            .await
2415            .expect("new_transaction failed");
2416        let handle = root_directory
2417            .create_child_file(&mut transaction, "test")
2418            .await
2419            .expect("create_child_file failed");
2420        transaction.commit().await.expect("commit failed");
2421
2422        let mut buf = handle.allocate_buffer(4096).await;
2423        buf.fill(0xaa);
2424        loop {
2425            if handle.write_or_append(None, buf.as_ref()).await.is_err() {
2426                break;
2427            }
2428        }
2429
2430        let max_offset = fs.allocator().maximum_offset();
2431        fs.close().await.expect("Close failed");
2432
2433        let device = fs.take_device().await;
2434        device.reopen(false);
2435        let mut buffer = device
2436            .allocate_buffer(
2437                crate::round::round_up(max_offset, TEST_DEVICE_BLOCK_SIZE).unwrap() as usize
2438            )
2439            .await;
2440        device.read(0, buffer.as_mut()).await.expect("read failed");
2441
2442        let image_data = buffer.to_vec();
2443        let device = DeviceHolder::new(
2444            FakeDevice::from_image(image_data.as_slice(), TEST_DEVICE_BLOCK_SIZE)
2445                .expect("from_image failed"),
2446        );
2447        let fs =
2448            FxFilesystemBuilder::new().read_only(true).open(device).await.expect("open failed");
2449        fs.close().await.expect("Close failed");
2450    }
2451
2452    #[test_case(SuperBlockInstance::A; "Superblock instance A")]
2453    #[test_case(SuperBlockInstance::B; "Superblock instance B")]
2454    #[fuchsia::test]
2455    async fn test_image_builder_mode_flush_on_close_sb_a(target_sb: SuperBlockInstance) {
2456        const BLOCK_SIZE: u32 = 4096;
2457        let device = DeviceHolder::new(FakeDevice::new(2048, BLOCK_SIZE));
2458
2459        // 1. Initialize in image_builder_mode
2460        device.reopen(true);
2461        let fs = FxFilesystemBuilder::new()
2462            .format(true)
2463            .image_builder_mode(Some(target_sb))
2464            .open(device)
2465            .await
2466            .expect("open failed");
2467
2468        fs.enable_allocations();
2469
2470        // 2. Finalize logic (via close)
2471        fs.device().reopen(false);
2472
2473        // 3. Write data
2474        {
2475            let root_store = fs.root_store();
2476            let root_directory =
2477                Directory::open(&root_store, root_store.root_directory_object_id())
2478                    .await
2479                    .expect("open failed");
2480
2481            let mut transaction = fs
2482                .root_store()
2483                .new_transaction(
2484                    lock_keys![LockKey::object(
2485                        root_directory.store().store_object_id(),
2486                        root_directory.object_id()
2487                    )],
2488                    Options::default(),
2489                )
2490                .await
2491                .expect("new transaction");
2492            let handle = root_directory
2493                .create_child_file(&mut transaction, "post_finalize_file")
2494                .await
2495                .expect("create file");
2496            transaction.commit().await.expect("commit");
2497
2498            let mut buf = handle.allocate_buffer(BLOCK_SIZE as usize).await;
2499            buf.fill(0xaa);
2500            handle.write_or_append(None, buf.as_ref()).await.expect("write failed");
2501        }
2502
2503        // 4. Close. Should flush to `target_sb` only.
2504        fs.close().await.expect("close failed");
2505
2506        let other_sb = target_sb.next();
2507
2508        // 5. Verify `target_sb` is valid and `other_sb` is empty.
2509        let device = fs.take_device().await;
2510        device.reopen(true); // Read-only is fine for verifying.
2511        let mut buf = device.allocate_buffer(BLOCK_SIZE as usize).await;
2512
2513        device.read(target_sb.first_extent().start, buf.as_mut()).await.expect("read target_sb");
2514        let data = buf.to_vec();
2515        assert_eq!(&data[..8], b"FxfsSupr", "target_sb should have magic bytes");
2516
2517        buf.fill(0); // Clear buffer
2518        device.read(other_sb.first_extent().start, buf.as_mut()).await.expect("read other_sb");
2519        // Expecting all zeros for `other_sb`
2520        let data2 = buf.to_vec();
2521        assert_eq!(data2, &[0; 4096], "other_sb should be zeroed");
2522    }
2523
2524    #[cfg(target_os = "fuchsia")]
2525    #[fuchsia::test(allow_stalls = false)]
2526    async fn test_trim_with_power_manager() {
2527        use anyhow::Error;
2528        use async_trait::async_trait;
2529        use fuchsia_async::TestExecutor;
2530        use futures::StreamExt;
2531
2532        TestExecutor::advance_to(fasync::MonotonicInstant::ZERO).await;
2533
2534        #[derive(Default)]
2535        struct MockPowerManager {
2536            on_battery: Mutex<bool>,
2537            event: event_listener::Event,
2538            wake_lease: Mutex<Option<zx::EventPair>>,
2539        }
2540
2541        impl MockPowerManager {
2542            fn set_on_battery(&self, v: bool) {
2543                *self.on_battery.lock() = v;
2544                self.event.notify(usize::MAX);
2545            }
2546
2547            fn is_lease_held(&self) -> bool {
2548                self.wake_lease.lock().as_ref().is_some_and(|handle| {
2549                    handle
2550                        .wait_one(
2551                            zx::Signals::EVENTPAIR_PEER_CLOSED,
2552                            zx::MonotonicInstant::INFINITE_PAST,
2553                        )
2554                        .is_err()
2555                })
2556            }
2557        }
2558
2559        impl super::PowerManager for MockPowerManager {
2560            fn watch_battery(
2561                self: Arc<Self>,
2562            ) -> futures::stream::BoxStream<'static, (bool, Option<super::WakeLease>)> {
2563                futures::stream::unfold(true, move |first| {
2564                    let this = self.clone();
2565                    async move {
2566                        if !first {
2567                            this.event.listen().await;
2568                        }
2569                        let val = *this.on_battery.lock();
2570                        let handle = if val {
2571                            None
2572                        } else {
2573                            let (h1, h2) = zx::EventPair::create();
2574                            *this.wake_lease.lock() = Some(h2);
2575                            Some(super::WakeLease::new(h1))
2576                        };
2577                        Some(((val, handle), false))
2578                    }
2579                })
2580                .boxed()
2581            }
2582        }
2583
2584        let trim_count = Arc::new(AtomicU32::new(0));
2585
2586        struct TrimTrackingDevice {
2587            inner: DeviceHolder,
2588            trim_count: Arc<AtomicU32>,
2589            power_manager: Arc<MockPowerManager>,
2590        }
2591
2592        #[async_trait]
2593        impl storage_device::Device for TrimTrackingDevice {
2594            fn allocate_buffer(&self, size: usize) -> storage_device::buffer::BufferFuture<'_> {
2595                self.inner.allocate_buffer(size)
2596            }
2597            fn block_size(&self) -> u32 {
2598                self.inner.block_size()
2599            }
2600            fn block_count(&self) -> u64 {
2601                self.inner.block_count()
2602            }
2603            async fn read_with_opts(
2604                &self,
2605                offset: u64,
2606                buffer: storage_device::buffer::MutableBufferRef<'_>,
2607                opts: storage_device::ReadOptions,
2608            ) -> Result<(), Error> {
2609                self.inner.read_with_opts(offset, buffer, opts).await
2610            }
2611            async fn write_with_opts(
2612                &self,
2613                offset: u64,
2614                buffer: storage_device::buffer::BufferRef<'_>,
2615                opts: storage_device::WriteOptions,
2616            ) -> Result<(), Error> {
2617                self.inner.write_with_opts(offset, buffer, opts).await
2618            }
2619            async fn trim(&self, range: std::ops::Range<u64>) -> Result<(), Error> {
2620                assert!(self.power_manager.is_lease_held());
2621                self.trim_count.fetch_add(1, Ordering::SeqCst);
2622                self.inner.trim(range).await
2623            }
2624            async fn flush(&self) -> Result<(), Error> {
2625                self.inner.flush().await
2626            }
2627            async fn close(&self) -> Result<(), Error> {
2628                self.inner.close().await
2629            }
2630            fn supports_trim(&self) -> bool {
2631                true
2632            }
2633            fn is_read_only(&self) -> bool {
2634                self.inner.is_read_only()
2635            }
2636            fn snapshot(&self) -> Result<DeviceHolder, Error> {
2637                Ok(DeviceHolder::new(TrimTrackingDevice {
2638                    inner: self.inner.snapshot()?,
2639                    trim_count: self.trim_count.clone(),
2640                    power_manager: self.power_manager.clone(),
2641                }))
2642            }
2643            fn reopen(&self, read_only: bool) {
2644                self.inner.reopen(read_only)
2645            }
2646        }
2647
2648        let pm = Arc::new(MockPowerManager::default());
2649
2650        // Start on battery.
2651        pm.set_on_battery(true);
2652
2653        let fake_device = FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE);
2654        let device = DeviceHolder::new(TrimTrackingDevice {
2655            inner: DeviceHolder::new(fake_device),
2656            trim_count: trim_count.clone(),
2657            power_manager: pm.clone(),
2658        });
2659
2660        let fs = FxFilesystemBuilder::new()
2661            .format(true)
2662            .power_manager(pm.clone())
2663            .trim_config(Some((Duration::ZERO, Duration::from_millis(100))))
2664            .trim_charger_wait(Duration::from_millis(10))
2665            .open(device)
2666            .await
2667            .expect("open failed");
2668
2669        // Initially on battery, so no trim should happen.
2670        TestExecutor::advance_to(fasync::MonotonicInstant::after(
2671            Duration::from_millis(500).into(),
2672        ))
2673        .await;
2674        let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2675
2676        assert_eq!(trim_count.load(Ordering::SeqCst), 0);
2677
2678        // Make some things to trim.
2679        {
2680            let root_store = fs.root_store();
2681            let root_directory =
2682                Directory::open(&root_store, root_store.root_directory_object_id())
2683                    .await
2684                    .expect("open failed");
2685            let mut transaction = fs
2686                .root_store()
2687                .new_transaction(
2688                    lock_keys![LockKey::object(
2689                        root_store.store_object_id(),
2690                        root_directory.object_id()
2691                    )],
2692                    Options::default(),
2693                )
2694                .await
2695                .expect("new_transaction failed");
2696            let handle = root_directory
2697                .create_child_file(&mut transaction, "test")
2698                .await
2699                .expect("create_child_file failed");
2700            transaction.commit().await.expect("commit failed");
2701            handle.allocate(0..4096).await.expect("allocate failed");
2702            // Now delete it to make it trimmable.
2703            let mut transaction = fs
2704                .root_store()
2705                .new_transaction(
2706                    lock_keys![
2707                        LockKey::object(root_store.store_object_id(), root_directory.object_id()),
2708                        LockKey::object(root_store.store_object_id(), handle.object_id()),
2709                    ],
2710                    Options::default(),
2711                )
2712                .await
2713                .expect("new_transaction failed");
2714            replace_child(&mut transaction, None, (&root_directory, "test"))
2715                .await
2716                .expect("delete failed");
2717            transaction.commit().await.expect("commit failed");
2718            fs.root_store()
2719                .tombstone_object(handle.object_id(), Options::default(), None)
2720                .await
2721                .expect("tombstone failed");
2722        }
2723
2724        // Put on external power source.
2725        pm.set_on_battery(false);
2726
2727        // Trim should start after 10ms.
2728        TestExecutor::advance_to(fasync::MonotonicInstant::after(Duration::from_millis(10).into()))
2729            .await;
2730
2731        let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2732
2733        assert!(trim_count.load(Ordering::SeqCst) > 0);
2734
2735        // Reset trim count and take off charger.
2736        trim_count.store(0, Ordering::SeqCst);
2737        pm.set_on_battery(true);
2738
2739        // Wait and ensure no more trims.
2740        TestExecutor::advance_to(fasync::MonotonicInstant::after(
2741            Duration::from_millis(500).into(),
2742        ))
2743        .await;
2744
2745        let _ = TestExecutor::poll_until_stalled(std::future::pending::<()>()).await;
2746
2747        assert_eq!(trim_count.load(Ordering::SeqCst), 0);
2748
2749        fs.close().await.expect("close failed");
2750    }
2751
2752    #[fuchsia::test]
2753    async fn test_concurrent_do_trim_returns_error() {
2754        let device = DeviceHolder::new(FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE));
2755        let fs = FxFilesystemBuilder::new()
2756            .trim_config(None)
2757            .format(true)
2758            .open(device)
2759            .await
2760            .expect("open failed");
2761
2762        let max_extent_size = fs.device().size() as usize;
2763        const EXTENTS_PER_BATCH: usize = usize::MAX;
2764
2765        // Hold onto trimmable extents to simulate an in-flight trim operation.
2766        let allocator = fs.allocator();
2767        let _trimmable_extents = allocator
2768            .take_for_trimming(0, max_extent_size, EXTENTS_PER_BATCH)
2769            .await
2770            .expect("take_for_trimming failed");
2771
2772        // Attempting to run do_trim concurrently while a trim is in-flight
2773        // should return FxfsError::AlreadyBound rather than panicking.
2774        let res = fs.do_trim(None).await;
2775        assert!(matches!(res, Err(e) if FxfsError::AlreadyBound.matches(&e)));
2776
2777        fs.close().await.expect("close failed");
2778    }
2779}